Files
Aislo/B03_FileInput/B03_FileInput_Repository.py
T
eomsangdonandClaude Opus 5 3572694f73 feat(B03/B04): LAS 없이 도엽등고선으로 3D 서피스를 만들어 설계한다
- B03: 'LAS 없이 설계' 토글 — LAS 필수카드 비활성화, las_free 플래그로
  업로드·완료 검증 면제, 계획노선 CSV를 WF1 분석 입력으로 사용
- B04: 신규 Engine_SheetSurface — 도엽_등고선.geojson을 노선 bbox+300m
  직사각형으로 절취, 배수유역 엔진의 등고선 정점구름·Delaunay TIN 보간을
  재사용해 dtm_sheet.npz(1m 격자, LAS DTM과 동일 형식) + 프리뷰 glb 생성.
  build_surface_sampler('sheet','dtm')로 종·횡단·배수 하류 계산 무수정 동작
- WF1: las_free면 run_sheet_surface_analysis로 분기, sheet/dtm 자동 확정.
  VWorld·도엽 확보 블록을 download_geodata()로 추출해 두 경로가 공유
- LAS가 있어도 도엽 서피스를 함께 생성·등록(참고용), B04에
  '도엽등고 3D 서피스' 별도 컨테이너로 표시

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-30 14:27:41 +09:00

396 lines
13 KiB
Python

"""B03 input_files 테이블의 aiomysql Raw SQL 접근."""
import json
from pathlib import PurePosixPath
from typing import Any
from uuid import UUID
import aiomysql
async def create_input_file(
connection: aiomysql.Connection,
*,
project_id: UUID,
file_type: str,
original_filename: str,
relative_path: str,
file_size_bytes: int,
upload_by: int | None,
crs_epsg: int | None,
metadata: dict[str, Any],
) -> int:
"""업로드 원본 파일 메타데이터를 저장하고 생성된 ID를 반환한다."""
normalized_path = PurePosixPath(relative_path)
if normalized_path.is_absolute() or ".." in normalized_path.parts:
raise ValueError("DB에는 프로젝트 루트 기준 상대 경로만 저장할 수 있습니다.")
if normalized_path.parts[:2] != ("B03_FileInput", "input"):
raise ValueError("입력 파일 경로는 B03_FileInput/input 아래여야 합니다.")
async with connection.cursor() as cursor:
await cursor.execute(
"""
INSERT INTO input_files (
project_id,
file_type,
original_filename,
raw_file_path,
file_size_mb,
upload_by,
crs_epsg,
metadata,
status
)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, 'UPLOADED')
""",
(
str(project_id),
file_type,
original_filename,
normalized_path.as_posix(),
file_size_bytes / (1024 * 1024),
upload_by,
crs_epsg,
json.dumps(metadata, ensure_ascii=False),
),
)
input_file_id = cursor.lastrowid
if not input_file_id:
raise RuntimeError("input_files 레코드 생성 결과에 ID가 없습니다.")
return int(input_file_id)
async def get_project_input_readiness(
connection: aiomysql.Connection,
project_id: UUID,
) -> tuple[set[str], int | None, int | None]:
"""업로드 파일 유형, 최신 포인트클라우드 입력 ID, 최신 계획노선 CSV 입력 ID를 반환한다."""
async with connection.cursor(aiomysql.DictCursor) as cursor:
await cursor.execute(
"""
SELECT id, LOWER(file_type) AS file_type
FROM input_files
WHERE project_id = %s AND status IN ('UPLOADED', 'PROCESSED')
ORDER BY id DESC
""",
(str(project_id),),
)
rows = await cursor.fetchall()
file_types = {str(row["file_type"]) for row in rows if row.get("file_type")}
point_cloud_id = next(
(int(row["id"]) for row in rows if str(row.get("file_type") or "") in {"las", "laz"}),
None,
)
# LAS 없는 설계(2026-08-30)의 WF1 입력 — 계획노선 CSV가 분석 원천이 된다.
route_csv_id = next(
(int(row["id"]) for row in rows if str(row.get("file_type") or "") == "csv"),
None,
)
return file_types, point_cloud_id, route_csv_id
async def get_project_storage_relative_path(
connection: aiomysql.Connection,
project_id: UUID,
) -> str:
"""프로젝트의 검증된 저장소 상대 경로를 조회한다."""
async with connection.cursor() as cursor:
await cursor.execute(
"""
SELECT storage_path
FROM projects
WHERE id = %s AND deleted_at IS NULL
""",
(str(project_id),),
)
row = await cursor.fetchone()
if not row or not row[0]:
raise LookupError("프로젝트 또는 프로젝트 저장 경로를 찾을 수 없습니다.")
normalized_path = PurePosixPath(str(row[0]).replace("\\", "/"))
if normalized_path.is_absolute() or ".." in normalized_path.parts:
raise ValueError("프로젝트 저장 경로는 안전한 상대 경로여야 합니다.")
return normalized_path.as_posix()
async def create_upload_session(
connection: aiomysql.Connection,
*,
session_id: str,
project_id: UUID,
original_filename: str,
file_size_bytes: int,
chunk_size_bytes: int,
total_chunks: int,
) -> None:
"""청크 업로드 세션을 생성하거나 동일 세션 ID를 갱신한다."""
async with connection.cursor() as cursor:
await cursor.execute(
"""
INSERT INTO upload_sessions (
id,
project_id,
original_filename,
file_size_bytes,
chunk_size_bytes,
total_chunks,
completed_chunks,
status,
created_at,
updated_at
)
VALUES (%s, %s, %s, %s, %s, %s, 0, 'in_progress', NOW(), NOW())
ON DUPLICATE KEY UPDATE
original_filename = VALUES(original_filename),
file_size_bytes = VALUES(file_size_bytes),
chunk_size_bytes = VALUES(chunk_size_bytes),
total_chunks = VALUES(total_chunks),
status = 'in_progress',
updated_at = NOW()
""",
(
session_id,
str(project_id),
original_filename,
file_size_bytes,
chunk_size_bytes,
total_chunks,
),
)
async def get_upload_session(
connection: aiomysql.Connection,
*,
project_id: UUID,
session_id: str,
) -> dict[str, Any]:
"""청크 업로드 세션을 조회한다."""
async with connection.cursor(aiomysql.DictCursor) as cursor:
await cursor.execute(
"""
SELECT
id,
project_id,
original_filename,
file_size_bytes,
chunk_size_bytes,
total_chunks,
completed_chunks,
status
FROM upload_sessions
WHERE id = %s AND project_id = %s
""",
(session_id, str(project_id)),
)
row = await cursor.fetchone()
if not row:
raise LookupError("업로드 세션을 찾을 수 없습니다.")
return dict(row)
async def upsert_upload_chunk(
connection: aiomysql.Connection,
*,
session_id: str,
chunk_index: int,
chunk_hash: str,
size_bytes: int,
stored_at: str,
) -> int:
"""청크 저장 정보를 기록하고 완료 청크 수를 반환한다."""
async with connection.cursor() as cursor:
await cursor.execute(
"""
INSERT INTO upload_chunks (
session_id,
chunk_index,
chunk_hash,
size_bytes,
stored_at,
completed_at
)
VALUES (%s, %s, %s, %s, %s, NOW())
ON DUPLICATE KEY UPDATE
chunk_hash = VALUES(chunk_hash),
size_bytes = VALUES(size_bytes),
stored_at = VALUES(stored_at),
completed_at = NOW()
""",
(session_id, chunk_index, chunk_hash, size_bytes, stored_at),
)
await cursor.execute(
"""
SELECT COUNT(*)
FROM upload_chunks
WHERE session_id = %s
""",
(session_id,),
)
row = await cursor.fetchone()
completed_chunks = int(row[0]) if row else 0
await cursor.execute(
"""
UPDATE upload_sessions
SET completed_chunks = %s, updated_at = NOW()
WHERE id = %s
""",
(completed_chunks, session_id),
)
return completed_chunks
async def list_completed_chunk_indexes(
connection: aiomysql.Connection,
*,
session_id: str,
) -> list[int]:
"""완료된 청크 인덱스 목록을 반환한다."""
async with connection.cursor() as cursor:
await cursor.execute(
"""
SELECT chunk_index
FROM upload_chunks
WHERE session_id = %s
ORDER BY chunk_index ASC
""",
(session_id,),
)
rows = await cursor.fetchall()
return [int(row[0]) for row in rows]
async def mark_upload_session_completed(
connection: aiomysql.Connection,
*,
session_id: str,
) -> None:
"""업로드 세션을 완료 상태로 표시한다."""
async with connection.cursor() as cursor:
await cursor.execute(
"""
UPDATE upload_sessions
SET status = 'completed', updated_at = NOW()
WHERE id = %s
""",
(session_id,),
)
async def mark_upload_session_failed(
connection: aiomysql.Connection,
*,
session_id: str,
) -> None:
"""업로드 세션을 실패 상태로 표시한다."""
async with connection.cursor() as cursor:
await cursor.execute(
"""
UPDATE upload_sessions
SET status = 'failed', updated_at = NOW()
WHERE id = %s
""",
(session_id,),
)
async def list_project_input_files(
connection: aiomysql.Connection,
project_id: UUID,
) -> list[dict[str, Any]]:
"""업로드 완료된 입력 파일 목록(재접속 현황 표시용) — 서버가 정본이다.
같은 파일명을 다시 올리면 새 레코드가 쌓이므로 파일명별 최신 것만 남긴다.
"""
async with connection.cursor(aiomysql.DictCursor) as cursor:
await cursor.execute(
"""
SELECT f.id, f.file_type, f.original_filename, f.file_size_mb, f.status,
f.upload_at
FROM input_files f
INNER JOIN (
SELECT MAX(id) AS id
FROM input_files
WHERE project_id = %s AND status IN ('UPLOADED', 'PROCESSED')
GROUP BY original_filename
) latest ON latest.id = f.id
ORDER BY f.id ASC
""",
(str(project_id),),
)
rows = await cursor.fetchall()
return [dict(row) for row in rows]
async def list_incomplete_upload_sessions(
connection: aiomysql.Connection,
project_id: UUID,
) -> list[dict[str, Any]]:
"""중단된(미완료) 청크 업로드 세션 목록 — 재접속 시 이어올리기 안내용."""
async with connection.cursor(aiomysql.DictCursor) as cursor:
await cursor.execute(
"""
SELECT id, original_filename, file_size_bytes, chunk_size_bytes,
total_chunks, completed_chunks, updated_at
FROM upload_sessions
WHERE project_id = %s AND status = 'in_progress'
ORDER BY updated_at DESC
""",
(str(project_id),),
)
rows = await cursor.fetchall()
return [dict(row) for row in rows]
async def find_input_file_by_name(
connection: aiomysql.Connection,
project_id: UUID,
original_filename: str,
) -> dict[str, Any] | None:
"""같은 이름으로 등록된 최신 입력 파일 1건. 없으면 None.
같은 파일을 다시 올렸는지 가리는 데 쓴다 — 중복 판정 기준은 파일명이고, 내용이 같은지는
이 행의 메타데이터에 적힌 지문으로 본다(2026-08-08 사용자 결정).
"""
async with connection.cursor(aiomysql.DictCursor) as cursor:
await cursor.execute(
"""
SELECT id, original_filename, file_size_mb, metadata
FROM input_files
WHERE project_id = %s AND original_filename = %s
AND status IN ('UPLOADED', 'PROCESSED')
ORDER BY id DESC
LIMIT 1
""",
(str(project_id), original_filename),
)
row = await cursor.fetchone()
return dict(row) if row else None
async def supersede_previous_input_files(
connection: aiomysql.Connection,
project_id: UUID,
original_filename: str,
keep_input_file_id: int,
) -> int:
"""같은 이름의 옛 행을 `SUPERSEDED`로 내린다. 내린 건수를 돌려준다.
조회 쿼리들이 `UPLOADED`/`PROCESSED`만 보므로, 이렇게만 해도 목록·분석에서 빠진다.
행을 지우지 않는 이유는 언제 무엇이 교체됐는지 추적할 근거를 남기기 위해서다.
"""
async with connection.cursor() as cursor:
await cursor.execute(
"""
UPDATE input_files
SET status = 'SUPERSEDED'
WHERE project_id = %s AND original_filename = %s AND id <> %s
AND status IN ('UPLOADED', 'PROCESSED')
""",
(str(project_id), original_filename, keep_input_file_id),
)
return cursor.rowcount