"""B04 지표면 분석 결과의 aiomysql Raw SQL 접근. processed_point_cloud(변환 포인트클라우드), surface_models(지표면 모델), terrain_layers(지형 레이어) 테이블에 메타데이터와 상대 경로를 기록한다. 공간 데이터는 MariaDB JSON 컬럼에 GeoJSON 문자열로 저장한다. """ import json from pathlib import Path, PurePosixPath from typing import Any from uuid import UUID import aiomysql _STAGE_ROOT = "B04_PreProcess" def _validate_stage_path(relative_path: str) -> str: """B04_PreProcess 아래의 안전한 상대 경로인지 검증하고 posix 문자열로 반환한다.""" normalized = PurePosixPath(relative_path.replace("\\", "/")) if normalized.is_absolute() or ".." in normalized.parts: raise ValueError("DB에는 프로젝트 루트 기준 상대 경로만 저장할 수 있습니다.") if not normalized.parts or normalized.parts[0] != _STAGE_ROOT: raise ValueError(f"B04 산출물 경로는 {_STAGE_ROOT} 아래여야 합니다.") return normalized.as_posix() async def create_processed_point_cloud( connection: aiomysql.Connection, *, input_file_id: int, project_id: UUID, process_type: str, processed_file_path: str | None, converted_format: str | None, converted_file_path: str | None, point_count: int | None, bounds: dict[str, Any] | None, statistics: dict[str, Any] | None, classification_summary: dict[str, Any] | None, processing_params: dict[str, Any] | None, status: str = "COMPLETE", ) -> int: """변환 포인트클라우드 메타데이터를 저장하고 생성된 ID를 반환한다.""" processed_rel = _validate_stage_path(processed_file_path) if processed_file_path else None converted_rel = _validate_stage_path(converted_file_path) if converted_file_path else None stats = statistics or {} async with connection.cursor() as cursor: await cursor.execute( """ INSERT INTO processed_point_cloud ( input_file_id, project_id, process_type, processed_file_path, converted_format, converted_file_path, point_count, min_z, max_z, mean_z, x_min, x_max, y_min, y_max, density_per_sqm, classification_summary, processing_params, status ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) """, ( input_file_id, str(project_id), process_type, processed_rel, converted_format, converted_rel, point_count, stats.get("min_z"), stats.get("max_z"), stats.get("mean_z"), (bounds or {}).get("x_min"), (bounds or {}).get("x_max"), (bounds or {}).get("y_min"), (bounds or {}).get("y_max"), stats.get("density_per_sqm"), json.dumps(classification_summary, ensure_ascii=False) if classification_summary is not None else None, json.dumps(processing_params, ensure_ascii=False) if processing_params is not None else None, status, ), ) new_id = cursor.lastrowid if not new_id: raise RuntimeError("processed_point_cloud 레코드 생성 결과에 ID가 없습니다.") return int(new_id) async def create_surface_model( connection: aiomysql.Connection, *, project_id: UUID, model_type: str, source_file_id: int | None, processed_cloud_id: int | None, crs_epsg: int | None, resolution_m: float | None, model_file_path: str | None, generation_params: dict[str, Any] | None, status: str = "COMPLETE", ) -> int: """지표면 모델 메타데이터를 저장하고 생성된 ID를 반환한다.""" model_rel = _validate_stage_path(model_file_path) if model_file_path else None async with connection.cursor() as cursor: await cursor.execute( """ INSERT INTO surface_models ( project_id, model_type, source_file_id, processed_cloud_id, status, crs_epsg, resolution_m, model_file_path, generation_params, completed_at ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, CURRENT_TIMESTAMP) """, ( str(project_id), model_type, source_file_id, processed_cloud_id, status, crs_epsg, resolution_m, model_rel, json.dumps(generation_params, ensure_ascii=False) if generation_params is not None else None, ), ) new_id = cursor.lastrowid if not new_id: raise RuntimeError("surface_models 레코드 생성 결과에 ID가 없습니다.") return int(new_id) async def create_terrain_layer( connection: aiomysql.Connection, *, surface_model_id: int, layer_name: str, geometry_type: str, layer_file_path: str | None, file_format: str | None, file_size_mb: float | None, statistics: dict[str, Any] | None, ) -> int: """지형 레이어 메타데이터를 저장하고 생성된 ID를 반환한다.""" layer_rel = _validate_stage_path(layer_file_path) if layer_file_path else None async with connection.cursor() as cursor: await cursor.execute( """ INSERT INTO terrain_layers ( surface_model_id, layer_name, geometry_type, layer_file_path, file_format, file_size_mb, statistics ) VALUES (%s, %s, %s, %s, %s, %s, %s) """, ( surface_model_id, layer_name, geometry_type, layer_rel, file_format, file_size_mb, json.dumps(statistics, ensure_ascii=False) if statistics is not None else None, ), ) new_id = cursor.lastrowid if not new_id: raise RuntimeError("terrain_layers 레코드 생성 결과에 ID가 없습니다.") return int(new_id) async def get_input_file( connection: aiomysql.Connection, project_id: UUID, input_file_id: int ) -> dict[str, Any]: """프로젝트의 특정 입력 파일 경로·좌표계를 조회한다.""" async with connection.cursor() as cursor: await cursor.execute( """ SELECT id, file_type, raw_file_path, crs_epsg FROM input_files WHERE id = %s AND project_id = %s """, (input_file_id, str(project_id)), ) row = await cursor.fetchone() if not row: raise LookupError("입력 파일을 찾을 수 없습니다.") return { "id": int(row[0]), "file_type": row[1], "raw_file_path": row[2], "crs_epsg": row[3], } async def list_project_point_cloud_inputs( connection: aiomysql.Connection, project_id: UUID ) -> list[dict[str, Any]]: """프로젝트의 LAS/LAZ 원본 입력 파일 목록을 최신순으로 조회한다.""" async with connection.cursor() as cursor: await cursor.execute( """ SELECT id, file_type, original_filename, raw_file_path, file_size_mb, crs_epsg, status, upload_at FROM input_files WHERE project_id = %s AND file_type IN ('las', 'laz') ORDER BY upload_at DESC, id DESC """, (str(project_id),), ) rows = await cursor.fetchall() return [ { "id": int(row[0]), "file_type": row[1], "original_filename": row[2], "raw_file_path": row[3], "file_size_mb": float(row[4]) if row[4] is not None else None, "crs_epsg": row[5], "status": row[6], "created_at": row[7].isoformat() if row[7] else None, } for row in rows ] async def list_project_point_cloud_paths( connection: aiomysql.Connection, project_id: UUID, project_root: Path ) -> list[Path]: """전처리가 병합할 지형 파일 경로 목록 — 실제로 있는 파일만 (2026-09-06 다중 입력). 교체된 옛 행(`SUPERSEDED`)은 조회에서 이미 빠지므로, 지금 살아 있는 파일만 남는다. """ rows = await list_project_point_cloud_inputs(connection, project_id) paths: list[Path] = [] for row in rows: raw = str(row.get("raw_file_path") or "") if not raw: continue path = project_root / Path(raw) if path.is_file() and path not in paths: paths.append(path) return paths async def list_surface_models( connection: aiomysql.Connection, project_id: UUID ) -> list[dict[str, Any]]: """프로젝트의 지표면 모델 목록을 최신순으로 조회한다.""" async with connection.cursor() as cursor: await cursor.execute( """ SELECT id, model_type, status, resolution_m, model_file_path, generation_params, created_at FROM surface_models WHERE project_id = %s ORDER BY created_at DESC """, (str(project_id),), ) rows = await cursor.fetchall() models: list[dict[str, Any]] = [] for row in rows: generation_params = row[5] if isinstance(generation_params, str): generation_params = json.loads(generation_params) models.append( { "id": int(row[0]), "model_type": row[1], "status": row[2], "resolution_m": row[3], "model_file_path": row[4], "generation_params": generation_params, "created_at": row[6].isoformat() if row[6] else None, } ) return models async def clear_confirmed_surface_models(connection: aiomysql.Connection, project_id: UUID) -> None: """재분석 시 기존 확정을 해제하여 새 분석 결과의 재확정을 요구한다.""" async with connection.cursor() as cursor: await cursor.execute( """ UPDATE surface_models SET status = 'COMPLETE' WHERE project_id = %s AND status = 'CONFIRMED' """, (str(project_id),), ) async def confirm_surface_model( connection: aiomysql.Connection, project_id: UUID, model_id: int ) -> None: """프로젝트 내 단일 모델만 CONFIRMED가 되도록 갱신한다.""" async with connection.cursor() as cursor: await cursor.execute( """ UPDATE surface_models SET status = 'COMPLETE' WHERE project_id = %s AND status = 'CONFIRMED' """, (str(project_id),), ) await cursor.execute( """ UPDATE surface_models SET status = 'CONFIRMED' WHERE id = %s AND project_id = %s AND status = 'COMPLETE' """, (model_id, str(project_id)), ) if cursor.rowcount != 1: raise LookupError("확정할 지표면 모델을 찾을 수 없습니다.") async def update_project_status( connection: aiomysql.Connection, project_id: UUID, status: str ) -> None: async with connection.cursor() as cursor: await cursor.execute( """ UPDATE projects SET status = %s, updated_at = NOW() WHERE id = %s AND deleted_at IS NULL """, (status, str(project_id)), ) async def delete_project_surface_models(connection: aiomysql.Connection, project_id: UUID) -> int: """프로젝트의 기존 지표면 모델 행을 모두 제거한다. 모든 분석 세대가 동일 파일 경로를 재사용하므로 재분석 시 구세대 행을 남기면 존재하지 않는 파일을 가리키게 된다. terrain_layers는 FK CASCADE, routes.surface_model_id는 FK SET NULL로 함께 정리된다. """ async with connection.cursor() as cursor: await cursor.execute( "DELETE FROM surface_models WHERE project_id = %s", (str(project_id),), ) return int(cursor.rowcount) async def save_surface_analysis_to_db( connection: aiomysql.Connection, *, project_id: UUID, input_file_id: int, analysis_result: dict[str, Any], source_filters: list[str], ) -> list[int]: """WF1 분석 결과를 DB에 저장한다 (기존 모델 행은 교체). 이 함수는 트랜잭션을 시작하거나 종료하지 않는다. 호출자는 같은 커넥션에서 begin/commit/rollback을 한 번만 수행해야 한다. """ input_file = await get_input_file(connection, project_id, input_file_id) await delete_project_surface_models(connection, project_id) processed = analysis_result["processed"] processed_cloud_id = await create_processed_point_cloud( connection, input_file_id=input_file_id, project_id=project_id, process_type="structured", processed_file_path=processed["processed_file_path"], converted_format=None, converted_file_path=processed["converted_file_path"], point_count=processed["point_count"], bounds=processed["bounds"], statistics=processed["statistics"], classification_summary=None, processing_params={"filters": source_filters}, ) surface_model_ids: list[int] = [] for model in analysis_result["models"]: model_id = await create_surface_model( connection, project_id=project_id, model_type=model["model_type"], source_file_id=input_file_id, processed_cloud_id=processed_cloud_id, crs_epsg=input_file["crs_epsg"], resolution_m=model["resolution_m"], model_file_path=model["model_file_path"], generation_params=model["generation_params"], ) surface_model_ids.append(model_id) for layer in model["layers"]: await create_terrain_layer( connection, surface_model_id=model_id, layer_name=layer["layer_name"], geometry_type=layer["geometry_type"], layer_file_path=layer["file_path"], file_format=layer["file_format"], file_size_mb=None, statistics=None, ) return surface_model_ids async def save_sheet_surface_models( connection: aiomysql.Connection, project_id: UUID, models: list[dict[str, Any]] ) -> list[int]: """관리자가 요청해 새로 만든 도엽 서피스 모델을 등록한다 (2026-09-01). 전체 재분석(`save_surface_analysis_to_db`)과 달리 **기존 행을 지우지 않는다** — 방식 하나만 덧붙이는 요청이라 이미 있는 모델을 날리면 안 된다. 같은 파일을 가리키는 옛 행만 걷어 낸다. 입력 파일·구조화 클라우드·좌표계는 같은 프로젝트의 기존 모델에서 물려받는다(도엽 서피스는 LAS가 아니라 등고선에서 나오지만 같은 사업지다). 이 함수는 트랜잭션을 시작하거나 끝내지 않는다 — 호출자가 한 번만 한다. """ async with connection.cursor() as cursor: await cursor.execute( """ SELECT source_file_id, processed_cloud_id, crs_epsg FROM surface_models WHERE project_id = %s ORDER BY id LIMIT 1 """, (str(project_id),), ) row = await cursor.fetchone() source_file_id, processed_cloud_id, crs_epsg = row if row else (None, None, None) model_ids: list[int] = [] for model in models: model_path = model.get("model_file_path") if model_path: async with connection.cursor() as cursor: await cursor.execute( "DELETE FROM surface_models WHERE project_id = %s AND model_file_path = %s", (str(project_id), model_path), ) model_id = await create_surface_model( connection, project_id=project_id, model_type=model["model_type"], source_file_id=source_file_id, processed_cloud_id=processed_cloud_id, crs_epsg=crs_epsg, resolution_m=model["resolution_m"], model_file_path=model_path, generation_params=model["generation_params"], ) model_ids.append(model_id) for layer in model.get("layers", []): await create_terrain_layer( connection, surface_model_id=model_id, layer_name=layer["layer_name"], geometry_type=layer["geometry_type"], layer_file_path=layer["file_path"], file_format=layer["file_format"], file_size_mb=None, statistics=None, ) return model_ids