Files
Aislo/B04_PreProcess/B04_PreProcess_Repository.py
T

450 lines
16 KiB
Python

"""B04 지표면 분석 결과의 aiomysql Raw SQL 접근.
processed_point_cloud(변환 포인트클라우드), surface_models(지표면 모델),
terrain_layers(지형 레이어) 테이블에 메타데이터와 상대 경로를 기록한다.
공간 데이터는 MariaDB JSON 컬럼에 GeoJSON 문자열로 저장한다.
"""
import json
from pathlib import 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_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