Files
Aislo/B03_FileInput/B03_FileInput_Router_Helpers.py
eomsangdonandClaude Opus 5 09e1c983ae feat(B03): 새 자료로 갈기 전에 「지우고 진행할까」 확인 한 단계 신설
업로드 하나가 그 프로젝트의 설계 산출물·초기값 스냅숏을 되돌릴 수 없게 지운다.
창 넷이 한 프로젝트를 볼 수 있어 말없이 지우면 남의 작업이 사라진다.

- `describe_existing_outputs`: 지워질 것을 사람 말로 낸다(전처리·노선·횡단·도면·
  수량·원가·초기값). 빈 폴더는 세지 않는다.
- 지울 것이 있는데 확인이 없으면 409 `confirm_required` — 무엇이 지워지는지 함께 준다.
- 세 갈래(직행·청크·보관함 연결) 전부에 `confirm_replace` 를 붙임.
- ⚠ 기본값은 확인 켬(True) — 자동 절차(체인·스크립트)는 물음에 안 걸린다.
  사람이 올리는 라우터만 False 로 불러 확인을 받는다.
- 화면: 409 를 받으면 「되돌릴 수 없습니다」와 지워질 목록을 보이고 한 번 묻는다.

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

364 lines
15 KiB
Python

"""B03 파일 입력 라우터 보조 — 필수 파일 판정·중복 업로드 판별·단계 기록·알림.
라우터 본체(`B03_FileInput_Router.py`)가 700줄을 넘어 떼어낸 조각이다(2026-09-04).
엔드포인트는 두지 않는다 — 순수 보조 함수와 파일 종류 상수만 둔다.
"""
import asyncio
import json
import logging
from pathlib import Path
from typing import Any
from uuid import UUID
import aiomysql
from fastapi.responses import JSONResponse
from B03_FileInput.B03_FileInput_Email import (
send_file_upload_complete_email,
)
from B03_FileInput.B03_FileInput_Repository import (
find_input_file_by_name,
get_project_input_readiness,
get_project_storage_relative_path,
)
from B03_FileInput.B03_FileInput_Schema import (
ChunkSessionCreateRequest,
ChunkSessionCreateResponse,
UploadedFileResult,
)
from common_util.common_util_initial_snapshot import clear_designing, discard_initial_snapshot
from common_util.common_util_json import atomic_write_json
from common_util.common_util_project_reset import purge_project_outputs
from common_util.common_util_storage import resolve_stored_project_path
from common_util.common_util_workflow_state import (
complete_stage,
reset_stages_after_input_change,
)
from config.config_db import get_db_pool
logger = logging.getLogger(__name__)
_ANALYSIS_RUNNING_MESSAGE = "이 프로젝트는 지금 분석 중입니다. 끝난 뒤에 새 자료를 올려 주세요."
# 새 자료를 받으면 옛 산출물을 지운다 — 지워질 것이 있으면 **먼저 묻는다**(2026-09-08).
# 창 넷이 한 프로젝트를 함께 볼 수 있어, 말없이 지우면 남이 며칠 만든 설계가 사라진다.
# 되돌릴 길이 없다(초기값 스냅숏까지 함께 버린다).
_OUTPUT_STAGE_LABELS: tuple[tuple[str, str], ...] = (
("B04_PreProcess", "지표면·전처리 결과"),
("B05_Profile", "노선·종단 설계"),
("B06_Section", "횡단 설계"),
("B07_DesignDetail", "도면"),
("B08_Quantity", "수량 산출"),
("B09_Estimation", "원가·내역"),
)
class OutputsWouldBeDiscarded(Exception):
"""이미 있는 설계 산출물을 지워야 새 자료를 받을 수 있는 상태.
라우터가 409 로 돌려주고, 사람이 「지우고 진행」을 고르면 `confirm_replace=True`
로 다시 들어온다. 자동 절차는 이 예외를 만들지 않는다 — 부를 때 확인을 켜 준다.
"""
def __init__(self, targets: list[str]) -> None:
self.targets = targets
super().__init__("설계 산출물이 있어 확인이 필요합니다.")
def _confirm_replace_response(exc: "OutputsWouldBeDiscarded") -> JSONResponse:
"""「지우고 진행할까?」를 사람에게 묻는 409 — 무엇이 지워지는지 함께 알린다."""
return JSONResponse(
status_code=409,
content={
"status": "confirm_required",
"confirm": "replace_outputs",
"message": (
"이 프로젝트에는 이미 만들어 둔 결과가 있습니다. 새 자료로 갈면 "
"아래가 지워지고 **되돌릴 수 없습니다**: " + " · ".join(exc.targets)
),
"targets": exc.targets,
},
)
def describe_existing_outputs(project_root: Path) -> list[str]:
"""새 자료를 받으면 **지워질 것**의 이름을 사람 말로 늘어놓는다. 없으면 빈 목록."""
targets: list[str] = []
for stage, label in _OUTPUT_STAGE_LABELS:
stage_root = project_root / stage
if not stage_root.is_dir():
continue
if any(entry.is_file() for entry in stage_root.rglob("*")):
targets.append(label)
if (project_root / "initial_snapshot").is_dir():
targets.append("초기값 스냅숏(되돌리기의 기준)")
return targets
_REQUIRED_FILE_TYPES = frozenset({"prj", "tfw"})
_POINT_CLOUD_FILE_TYPES = frozenset({"las", "laz"})
# 계획노선은 CSV 또는 shapefile 중 하나면 된다 (2026-08-31 — 원청 정식 노선이 shapefile).
_ROUTE_FILE_TYPES = frozenset({"csv", "shp"})
# shapefile은 이것들이 다 있어야 열린다. `.cpg`는 없으면 CP949로 읽으므로 필수가 아니다.
# `route_prj`는 노선 세트 폴더에 있는 PRJ — 지형 PRJ(`prj`)와 따로 센다.
_SHAPEFILE_REQUIRED_TYPES = frozenset({"shx", "dbf", "route_prj"})
def _total_chunks(size_bytes: int, chunk_size_bytes: int) -> int:
return max(1, (size_bytes + chunk_size_bytes - 1) // chunk_size_bytes)
def _is_point_cloud_result(result: UploadedFileResult) -> bool:
"""포인트클라우드 결과인지 — 임시 보관함 안내 메일 경로에서 쓴다.
프로젝트 업로드 경로는 더 이상 이 판정으로 메일을 보내지 않는다
([[_send_upload_complete_notification]] 주석 참고).
"""
return result.file_type.lower() in _POINT_CLOUD_FILE_TYPES
def upload_file_types(filenames: list[str]) -> set[str]:
"""업로드 요청의 파일명을 **완료 검사와 같은 기준**으로 유형화한다.
확장자만 세면 노선 세트의 `.prj`가 지형 PRJ로 오인돼, 요청 검사는 통과하고 저장 뒤
완료 검사에서 누락으로 갈린다(2026-09-06 실측 — LAS 없이 설계가 이 어긋남으로 막혔다).
노선 도형(`.shp`)과 basename이 같은 PRJ는 화면의 카드 배정과 같은 규칙으로 `route_prj`.
"""
route_stem = next(
(Path(name).stem for name in filenames if Path(name).suffix.lower() == ".shp"), None
)
types: set[str] = set()
for name in filenames:
suffix = Path(name).suffix.lower().lstrip(".")
if suffix == "prj" and route_stem is not None and Path(name).stem == route_stem:
suffix = "route_prj"
types.add(suffix)
return types
def _missing_required_file_types(file_types: set[str], las_free: bool = False) -> list[str]:
# LAS 없는 설계(도엽등고선 기반, 2026-08-30)는 지형 한 벌(포인트클라우드·지형 PRJ·
# TFW)을 통째로 받지 않는다 — 화면도 그 카드들을 필수에서 뺀다.
missing = [] if las_free else sorted(_REQUIRED_FILE_TYPES - file_types)
if not file_types.intersection(_ROUTE_FILE_TYPES):
missing.append("csv/shp")
# shapefile로 왔으면 형제 파일이 다 있어야 노선을 읽는다.
if "shp" in file_types:
missing.extend(sorted(_SHAPEFILE_REQUIRED_TYPES - file_types))
if not las_free and not file_types.intersection(_POINT_CLOUD_FILE_TYPES):
missing.append("las/laz")
return missing
def _require_complete_file_set(file_types: set[str], las_free: bool = False) -> None:
missing = _missing_required_file_types(file_types, las_free)
if missing:
raise ValueError(f"B03 필수 입력 파일이 없습니다: {', '.join(missing)}")
def _stored_fingerprint(metadata: Any) -> str | None:
"""입력 파일 메타데이터에 적어 둔 지문을 꺼낸다."""
if isinstance(metadata, str):
try:
metadata = json.loads(metadata)
except (TypeError, ValueError):
return None
if not isinstance(metadata, dict):
return None
value = metadata.get("fingerprint")
return str(value) if value else None
async def _already_uploaded(
connection: aiomysql.Connection,
project_id: UUID,
payload: ChunkSessionCreateRequest,
) -> ChunkSessionCreateResponse | None:
"""같은 이름으로 **같은 내용**이 이미 올라와 있으면 전송을 건너뛰라는 응답을 만든다.
1.7GB를 다 받은 뒤에 비교하면 아낄 게 없으므로, 세션을 만들기 전에 화면이 보내 준
지문으로 가린다. 지문이 없거나 다르면 그냥 올린다 — 애매하면 올리는 쪽이 안전하다.
"""
if not payload.fingerprint:
return None
existing = await find_input_file_by_name(connection, project_id, payload.original_filename)
if not existing or _stored_fingerprint(existing.get("metadata")) != payload.fingerprint:
return None
logger.info(
"B03 같은 파일 재업로드 — 전송 생략: project_id=%s file=%s",
project_id,
payload.original_filename,
)
return ChunkSessionCreateResponse(
project_id=str(project_id),
upload_session_id="",
original_filename=payload.original_filename,
file_size_bytes=payload.size_bytes,
chunk_size_bytes=payload.chunk_size_bytes,
total_chunks=0,
already_uploaded=True,
)
async def _complete_file_input_if_ready(
connection: aiomysql.Connection,
project_id: UUID,
las_free: bool = False,
confirm_replace: bool = True,
) -> int:
"""새 자료로 갈아 끼운다. `confirm_replace=False` 면 지울 것이 있을 때 멈추고 묻는다.
기본이 `True` 인 까닭 — 자동 절차(체인·스크립트)는 물음에 걸리면 안 된다.
사람이 올리는 갈래(라우터)만 `False` 로 불러 확인을 받는다(2026-09-08 사용자 지시).
"""
file_types, point_cloud_input_id, route_csv_input_id = await get_project_input_readiness(
connection, project_id
)
_require_complete_file_set(file_types, las_free)
if point_cloud_input_id is None:
if not las_free:
raise ValueError("LAS 또는 LAZ 입력 파일을 찾을 수 없습니다.")
if route_csv_input_id is None:
raise ValueError("계획 노선 입력 파일(CSV 또는 shapefile)을 찾을 수 없습니다.")
# 자료가 갈렸으니 옛 계산 결과(파일 + DB)를 지우고 진행 표시도 되돌린다. 남겨 두면
# 아직 다시 만들어지지 않은 뒷단계 화면이 옛 결과를 새 자료 것인 양 보여준다.
stored_path = await get_project_storage_relative_path(connection, project_id)
project_root = Path(resolve_stored_project_path(stored_path))
# 지형 파일 여러 장은 합쳐서 전처리한다 — 다른 사업지 파일이 섞이면 합친 범위가
# 통째로 어긋나므로 여기서 막는다(2026-09-06 사용자 지시). 머리글만 읽어 즉시 끝난다.
from B04_PreProcess.B04_PreProcess_Engine_Structurize import merge_gap_error
from B04_PreProcess.B04_PreProcess_Repository import list_project_point_cloud_paths
terrain_paths = await list_project_point_cloud_paths(connection, project_id, project_root)
gap_message = merge_gap_error(terrain_paths)
if gap_message:
raise ValueError(gap_message)
# ⚠ 여기서부터는 되돌릴 수 없다 — 아래 셋이 파일까지 지운다. 지울 것이 있으면 먼저 묻는다.
if not confirm_replace:
targets = describe_existing_outputs(project_root)
if targets:
raise OutputsWouldBeDiscarded(targets)
clear_designing(project_root)
discard_initial_snapshot(project_root)
await purge_project_outputs(connection, str(project_id), project_root)
async with connection.cursor(aiomysql.DictCursor) as cursor:
await reset_stages_after_input_change(cursor, str(project_id))
await complete_stage(cursor, str(project_id), 0)
# LAS가 있으면 LAS, 없으면(las_free) 계획노선 CSV가 WF1 분석 입력이다.
return point_cloud_input_id if point_cloud_input_id is not None else int(route_csv_input_id)
def _write_stage_metadata(
stage_root: Path,
project_id: UUID,
results: list[UploadedFileResult],
) -> None:
metadata_path = stage_root / "metadata.json"
existing_files: list[dict[str, Any]] = []
if metadata_path.exists():
try:
payload = json.loads(metadata_path.read_text(encoding="utf-8"))
existing_files = list(payload.get("files") or [])
except (OSError, TypeError, ValueError):
logger.warning("B03 metadata.json을 읽지 못해 새로 작성합니다: %s", metadata_path)
merged = {
str(item.get("relative_path") or item.get("original_filename")): item
for item in existing_files
}
for result in results:
dumped = result.model_dump()
merged[result.relative_path] = dumped
atomic_write_json(
metadata_path,
{"project_id": str(project_id), "files": list(merged.values())},
)
# 실행 중인 백그라운드 작업의 강한 참조. 이벤트 루프는 작업을 약한 참조로만 들고 있어,
# 여기서 붙잡지 않으면 GC가 대기 중인 작업을 통째로 회수해 WF1 분석이 조용히 사라진다
# (업로드는 성공·stage 0은 COMPLETE인데 stage 1은 NOT_STARTED로 남는 증상).
_BACKGROUND_TASKS: set[asyncio.Task] = set()
def _schedule_background_task(coro: Any, *, task_name: str) -> None:
task = asyncio.create_task(coro, name=task_name)
_BACKGROUND_TASKS.add(task)
def _log_task_failure(completed: asyncio.Task) -> None:
_BACKGROUND_TASKS.discard(completed)
try:
completed.result()
except asyncio.CancelledError:
logger.warning("백그라운드 작업 취소됨: %s", task_name)
except Exception:
logger.exception("백그라운드 작업 실패: %s", task_name)
task.add_done_callback(_log_task_failure)
logger.info("백그라운드 작업 시작: %s", task_name)
async def _get_project_notification_info(
connection: aiomysql.Connection,
project_id: UUID,
) -> dict[str, Any] | None:
async with connection.cursor(aiomysql.DictCursor) as cursor:
await cursor.execute(
"""
SELECT
p.id,
p.name AS project_name,
u.email AS user_email,
u.name AS user_name
FROM projects p
JOIN users u ON u.id = p.user_id
WHERE p.id = %s AND p.deleted_at IS NULL AND u.deleted_at IS NULL
""",
(str(project_id),),
)
row = await cursor.fetchone()
return dict(row) if row else None
async def _update_project_status(project_id: UUID, status: str) -> None:
pool = get_db_pool()
async with pool.acquire() as connection, 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)),
)
await connection.commit()
async def _send_upload_complete_notification(
*,
project_id: UUID,
uploaded_file: UploadedFileResult,
) -> None:
"""저장만 끝났을 때 보내는 안내.
프로젝트 업로드 경로에서는 호출하지 않는다 — 그 흐름은 초기 설계까지 마친 뒤
통합 메일 한 통으로 알린다. 프로젝트 생성 전 임시 보관함 업로드에서 쓸 예정이라
지워두지 않았다(2026-08-08 사용자 지시).
"""
pool = get_db_pool()
async with pool.acquire() as connection:
project_info = await _get_project_notification_info(connection, project_id)
if not project_info or not project_info.get("user_email"):
logger.warning("업로드 완료 이메일 수신자 없음: project_id=%s", project_id)
return
await send_file_upload_complete_email(
to_email=str(project_info["user_email"]),
user_name=str(project_info.get("user_name") or "사용자"),
project_name=str(project_info.get("project_name") or project_id),
file_name=uploaded_file.original_filename,
file_size_mb=uploaded_file.size_bytes / (1024 * 1024),
metadata=uploaded_file.metadata,
)