Files
Aislo/B03_FileInput/B03_FileInput_Router.py
T
eomsangdonandClaude Opus 5 ded617129e fix(B03): upload-overview 컬럼명 upload_at으로 수정 (input_files 스키마 정합)
서버 테스트에서 발견 — input_files에는 created_at이 없고 upload_at이다.
upload-overview 응답·Repository 조회 컬럼을 스키마에 맞춤.

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

694 lines
28 KiB
Python

"""B03 파일 입력 FastAPI 라우터."""
import asyncio
import json
import logging
from pathlib import Path
from typing import Any
from uuid import UUID, uuid4
import aiomysql
from fastapi import APIRouter, Depends, File, Form, UploadFile
from fastapi.responses import JSONResponse
from B03_FileInput.B03_FileInput_Email import (
send_file_upload_complete_email,
)
from B03_FileInput.B03_FileInput_Engine import (
merge_upload_chunks,
remove_chunk_session,
resolve_chunk_session_dir,
resolve_upload_destination,
save_upload_chunk,
save_upload_stream,
)
from B03_FileInput.B03_FileInput_Engine_Analyze import analyze_input_metadata
from B03_FileInput.B03_FileInput_Repository import (
create_input_file,
create_upload_session,
get_project_input_readiness,
get_project_storage_relative_path,
get_upload_session,
list_completed_chunk_indexes,
list_incomplete_upload_sessions,
list_project_input_files,
mark_upload_session_completed,
mark_upload_session_failed,
upsert_upload_chunk,
)
from B03_FileInput.B03_FileInput_Schema import (
ChunkSessionCreateRequest,
ChunkSessionCreateResponse,
ChunkUploadResponse,
FileUploadDescriptor,
FileUploadResponse,
UploadedFileResult,
UploadFinalizeRequest,
UploadOverviewFile,
UploadOverviewResponse,
UploadOverviewSession,
UploadStatusResponse,
)
from B03_FileInput.B03_FileInput_Service_WF1 import trigger_wf1_analysis_and_email
from common_util.common_util_auth import verify_session
from common_util.common_util_json import atomic_write_json
from common_util.common_util_storage import resolve_stored_project_path
from common_util.common_util_workflow import load_project_workflow
from common_util.common_util_workflow_state import (
complete_stage,
get_workflow_state,
)
from config.config_db import get_db_pool
from config.config_system import (
UPLOAD_CHUNK_SIZE_BYTES,
UPLOAD_MAX_FILES,
)
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/projects", tags=["B03 File Input"])
_REQUIRED_FILE_TYPES = frozenset({"csv", "prj", "tfw"})
_POINT_CLOUD_FILE_TYPES = frozenset({"las", "laz"})
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:
return result.file_type.lower() in _POINT_CLOUD_FILE_TYPES
def _missing_required_file_types(file_types: set[str]) -> list[str]:
missing = sorted(_REQUIRED_FILE_TYPES - file_types)
if not file_types.intersection(_POINT_CLOUD_FILE_TYPES):
missing.append("las/laz")
return missing
def _require_complete_file_set(file_types: set[str]) -> None:
missing = _missing_required_file_types(file_types)
if missing:
raise ValueError(f"B03 필수 입력 파일이 없습니다: {', '.join(missing)}")
async def _complete_file_input_if_ready(
connection: aiomysql.Connection,
project_id: UUID,
) -> int:
file_types, point_cloud_input_id = await get_project_input_readiness(connection, project_id)
_require_complete_file_set(file_types)
if point_cloud_input_id is None:
raise ValueError("LAS 또는 LAZ 입력 파일을 찾을 수 없습니다.")
async with connection.cursor() as cursor:
await complete_stage(cursor, str(project_id), 0)
return point_cloud_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())},
)
def _schedule_background_task(coro: Any, *, task_name: str) -> None:
task = asyncio.create_task(coro, name=task_name)
def _log_task_failure(completed: asyncio.Task) -> None:
try:
completed.result()
except Exception:
logger.exception("백그라운드 작업 실패: %s", task_name)
task.add_done_callback(_log_task_failure)
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:
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,
)
@router.post("/{project_id}/files", response_model=FileUploadResponse)
async def upload_project_files(
project_id: UUID,
files: list[UploadFile] = File(...),
session: dict[str, Any] = Depends(verify_session),
) -> FileUploadResponse | JSONResponse:
"""프로젝트 입력 파일을 저장·분석하고 DB 메타데이터를 기록한다."""
if not files or len(files) > UPLOAD_MAX_FILES:
return JSONResponse(
status_code=400,
content={
"status": "error",
"message": f"파일은 1~{UPLOAD_MAX_FILES}개까지 가능합니다.",
},
)
filenames = [upload.filename or "" for upload in files]
if len({filename.casefold() for filename in filenames}) != len(filenames):
return JSONResponse(
status_code=400,
content={"status": "error", "message": "동일한 파일명을 중복 업로드할 수 없습니다."},
)
las_count = sum(Path(filename).suffix.lower() in {".las", ".laz"} for filename in filenames)
if las_count != 1:
return JSONResponse(
status_code=400,
content={
"status": "error",
"message": "LAS 또는 LAZ 파일을 정확히 1개 포함해야 합니다.",
},
)
csv_count = sum(Path(filename).suffix.lower() == ".csv" for filename in filenames)
if csv_count != 1:
return JSONResponse(
status_code=400,
content={
"status": "error",
"message": "계획노선 CSV 파일을 정확히 1개 포함해야 합니다.",
},
)
request_file_types = {Path(filename).suffix.lower().lstrip(".") for filename in filenames}
missing_required = _missing_required_file_types(request_file_types)
if missing_required:
return JSONResponse(
status_code=400,
content={
"status": "error",
"message": f"B03 필수 입력 파일이 없습니다: {', '.join(missing_required)}",
},
)
pool = get_db_pool()
saved_paths: list[Path] = []
try:
results: list[UploadedFileResult] = []
point_cloud_input_id: int | None = None
async with pool.acquire() as connection:
stored_path = await get_project_storage_relative_path(connection, project_id)
project_root = Path(resolve_stored_project_path(stored_path))
await connection.begin()
try:
for upload in files:
preliminary = FileUploadDescriptor(
original_filename=upload.filename or "",
size_bytes=max(upload.size or 0, 1),
)
destination = resolve_upload_destination(project_root, preliminary)
written_bytes = await save_upload_stream(upload, destination)
saved_paths.append(destination)
descriptor = FileUploadDescriptor(
original_filename=preliminary.original_filename,
size_bytes=written_bytes,
)
metadata = await asyncio.to_thread(analyze_input_metadata, destination)
relative_path = destination.relative_to(project_root).as_posix()
file_type = destination.suffix.lower().lstrip(".")
crs_epsg = metadata.get("epsg")
input_file_id = await create_input_file(
connection,
project_id=project_id,
file_type=file_type,
original_filename=descriptor.original_filename,
relative_path=relative_path,
file_size_bytes=written_bytes,
upload_by=None,
crs_epsg=int(crs_epsg) if crs_epsg is not None else None,
metadata=metadata,
)
results.append(
UploadedFileResult(
input_file_id=input_file_id,
original_filename=descriptor.original_filename,
file_type=file_type,
relative_path=relative_path,
size_bytes=written_bytes,
metadata=metadata,
)
)
point_cloud_input_id = await _complete_file_input_if_ready(connection, project_id)
await connection.commit()
except Exception:
await connection.rollback()
raise
stage_root = project_root / "B03_FileInput"
_write_stage_metadata(stage_root, project_id, results)
workflow_path = project_root / "workflow.json"
if not workflow_path.exists():
atomic_write_json(workflow_path, load_project_workflow(project_root))
point_cloud_result = next(
(result for result in results if _is_point_cloud_result(result)),
None,
)
if point_cloud_result:
_schedule_background_task(
_send_upload_complete_notification(
project_id=project_id,
uploaded_file=point_cloud_result,
),
task_name=f"b03-upload-email-{project_id}",
)
if point_cloud_input_id is not None:
_schedule_background_task(
trigger_wf1_analysis_and_email(
project_id=project_id,
input_file_id=point_cloud_input_id,
user_role=str(session["role"]),
),
task_name=f"b04-wf1-auto-{project_id}",
)
return FileUploadResponse(project_id=str(project_id), files=results)
except LookupError as exc:
return JSONResponse(status_code=404, content={"status": "error", "message": str(exc)})
except (OSError, ValueError) as exc:
for saved_path in saved_paths:
saved_path.unlink(missing_ok=True)
return JSONResponse(status_code=400, content={"status": "error", "message": str(exc)})
except Exception:
for saved_path in saved_paths:
saved_path.unlink(missing_ok=True)
logger.exception("B03 파일 업로드 처리 실패: project_id=%s", project_id)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "파일 업로드 처리 중 오류가 발생했습니다."},
)
finally:
for upload in files:
await upload.close()
@router.post("/{project_id}/upload-sessions", response_model=ChunkSessionCreateResponse)
async def create_project_upload_session(
project_id: UUID,
payload: ChunkSessionCreateRequest,
) -> ChunkSessionCreateResponse | JSONResponse:
"""대용량 파일 청크 업로드 세션을 생성한다."""
chunk_size_bytes = min(payload.chunk_size_bytes, UPLOAD_CHUNK_SIZE_BYTES)
total_chunks = _total_chunks(payload.size_bytes, chunk_size_bytes)
session_id = str(uuid4())
pool = get_db_pool()
try:
async with pool.acquire() as connection:
await get_project_storage_relative_path(connection, project_id)
await create_upload_session(
connection,
session_id=session_id,
project_id=project_id,
original_filename=payload.original_filename,
file_size_bytes=payload.size_bytes,
chunk_size_bytes=chunk_size_bytes,
total_chunks=total_chunks,
)
return ChunkSessionCreateResponse(
project_id=str(project_id),
upload_session_id=session_id,
original_filename=payload.original_filename,
file_size_bytes=payload.size_bytes,
chunk_size_bytes=chunk_size_bytes,
total_chunks=total_chunks,
)
except LookupError as exc:
return JSONResponse(status_code=404, content={"status": "error", "message": str(exc)})
except (OSError, ValueError) as exc:
return JSONResponse(status_code=400, content={"status": "error", "message": str(exc)})
except Exception:
logger.exception("B03 청크 세션 생성 실패: project_id=%s", project_id)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "업로드 세션 생성 중 오류가 발생했습니다."},
)
@router.post("/{project_id}/chunks", response_model=ChunkUploadResponse)
async def upload_project_chunk(
project_id: UUID,
session_id: str = Form(...),
chunk_index: int = Form(...),
chunk_data: UploadFile = File(...),
) -> ChunkUploadResponse | JSONResponse:
"""단일 파일 청크를 B03 임시 폴더에 저장하고 DB에 기록한다."""
pool = get_db_pool()
try:
async with pool.acquire() as connection:
session = await get_upload_session(
connection,
project_id=project_id,
session_id=session_id,
)
if chunk_index < 0 or chunk_index >= int(session["total_chunks"]):
return JSONResponse(
status_code=400,
content={"status": "error", "message": "청크 인덱스가 범위를 벗어났습니다."},
)
stored_path = await get_project_storage_relative_path(connection, project_id)
project_root = Path(resolve_stored_project_path(stored_path))
session_dir = resolve_chunk_session_dir(project_root, session_id)
chunk_path, size_bytes, chunk_hash = await save_upload_chunk(
chunk_data,
session_dir,
chunk_index,
expected_max_bytes=int(session["chunk_size_bytes"]),
)
relative_chunk_path = chunk_path.relative_to(project_root).as_posix()
completed_chunks = await upsert_upload_chunk(
connection,
session_id=session_id,
chunk_index=chunk_index,
chunk_hash=chunk_hash,
size_bytes=size_bytes,
stored_at=relative_chunk_path,
)
return ChunkUploadResponse(
upload_session_id=session_id,
chunk_index=chunk_index,
completed_chunks=completed_chunks,
total_chunks=int(session["total_chunks"]),
chunk_hash=chunk_hash,
)
except LookupError as exc:
return JSONResponse(status_code=404, content={"status": "error", "message": str(exc)})
except (OSError, ValueError) as exc:
return JSONResponse(status_code=400, content={"status": "error", "message": str(exc)})
except Exception:
logger.exception(
"B03 청크 업로드 실패: project_id=%s session_id=%s",
project_id,
session_id,
)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "청크 업로드 중 오류가 발생했습니다."},
)
finally:
await chunk_data.close()
@router.post("/{project_id}/finalize", response_model=FileUploadResponse)
async def finalize_project_upload(
project_id: UUID,
payload: UploadFinalizeRequest,
session: dict[str, Any] = Depends(verify_session),
) -> FileUploadResponse | JSONResponse:
"""청크 업로드를 최종 병합하고 input_files 메타데이터를 기록한다."""
pool = get_db_pool()
final_path: Path | None = None
point_cloud_input_id: int | None = None
try:
async with pool.acquire() as connection:
session = await get_upload_session(
connection,
project_id=project_id,
session_id=payload.session_id,
)
if int(session["total_chunks"]) != payload.total_chunks:
return JSONResponse(
status_code=400,
content={"status": "error", "message": "세션 청크 개수가 일치하지 않습니다."},
)
completed_indexes = await list_completed_chunk_indexes(
connection,
session_id=payload.session_id,
)
expected_indexes = list(range(payload.total_chunks))
if completed_indexes != expected_indexes:
return JSONResponse(
status_code=400,
content={"status": "error", "message": "아직 업로드되지 않은 청크가 있습니다."},
)
stored_path = await get_project_storage_relative_path(connection, project_id)
project_root = Path(resolve_stored_project_path(stored_path))
descriptor = FileUploadDescriptor(
original_filename=str(session["original_filename"]),
size_bytes=int(session["file_size_bytes"]),
)
final_path = merge_upload_chunks(
project_root,
descriptor,
payload.session_id,
payload.total_chunks,
)
metadata = await asyncio.to_thread(analyze_input_metadata, final_path)
relative_path = final_path.relative_to(project_root).as_posix()
file_type = final_path.suffix.lower().lstrip(".")
crs_epsg = metadata.get("epsg")
await connection.begin()
try:
input_file_id = await create_input_file(
connection,
project_id=project_id,
file_type=file_type,
original_filename=descriptor.original_filename,
relative_path=relative_path,
file_size_bytes=descriptor.size_bytes,
upload_by=None,
crs_epsg=int(crs_epsg) if crs_epsg is not None else None,
metadata=metadata,
)
await mark_upload_session_completed(connection, session_id=payload.session_id)
if payload.complete_upload:
point_cloud_input_id = await _complete_file_input_if_ready(
connection,
project_id,
)
await connection.commit()
except Exception:
await connection.rollback()
await mark_upload_session_failed(connection, session_id=payload.session_id)
raise
remove_chunk_session(project_root, payload.session_id)
result = UploadedFileResult(
input_file_id=input_file_id,
original_filename=descriptor.original_filename,
file_type=file_type,
relative_path=relative_path,
size_bytes=descriptor.size_bytes,
metadata=metadata,
)
stage_root = project_root / "B03_FileInput"
_write_stage_metadata(stage_root, project_id, [result])
if _is_point_cloud_result(result):
_schedule_background_task(
_send_upload_complete_notification(
project_id=project_id,
uploaded_file=result,
),
task_name=f"b03-upload-email-{project_id}",
)
if point_cloud_input_id is not None:
_schedule_background_task(
trigger_wf1_analysis_and_email(
project_id=project_id,
input_file_id=point_cloud_input_id,
user_role=str(session["role"]),
),
task_name=f"b04-wf1-auto-{project_id}",
)
return FileUploadResponse(project_id=str(project_id), files=[result])
except LookupError as exc:
return JSONResponse(status_code=404, content={"status": "error", "message": str(exc)})
except (OSError, ValueError) as exc:
if final_path is not None:
final_path.unlink(missing_ok=True)
return JSONResponse(status_code=400, content={"status": "error", "message": str(exc)})
except Exception:
if final_path is not None:
final_path.unlink(missing_ok=True)
logger.exception("B03 청크 업로드 최종 병합 실패: project_id=%s", project_id)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "업로드 최종 처리 중 오류가 발생했습니다."},
)
@router.get("/{project_id}/upload-status/{session_id}", response_model=UploadStatusResponse)
async def get_project_upload_status(
project_id: UUID,
session_id: str,
) -> UploadStatusResponse | JSONResponse:
"""업로드 세션의 청크 완료 상태를 조회한다."""
pool = get_db_pool()
try:
async with pool.acquire() as connection:
session = await get_upload_session(
connection,
project_id=project_id,
session_id=session_id,
)
completed_indexes = await list_completed_chunk_indexes(
connection,
session_id=session_id,
)
return UploadStatusResponse(
upload_session_id=session_id,
upload_status=str(session["status"]),
original_filename=str(session["original_filename"]),
file_size_bytes=int(session["file_size_bytes"]),
chunk_size_bytes=int(session["chunk_size_bytes"]),
total_chunks=int(session["total_chunks"]),
completed_chunks=len(completed_indexes),
completed_chunk_indexes=completed_indexes,
)
except LookupError as exc:
return JSONResponse(status_code=404, content={"status": "error", "message": str(exc)})
except (OSError, ValueError) as exc:
return JSONResponse(status_code=400, content={"status": "error", "message": str(exc)})
except Exception:
logger.exception(
"B03 업로드 상태 조회 실패: project_id=%s session_id=%s",
project_id,
session_id,
)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "업로드 상태 조회 중 오류가 발생했습니다."},
)
@router.get("/{project_id}/upload-overview", response_model=UploadOverviewResponse)
async def get_project_upload_overview(
project_id: UUID,
session: dict[str, Any] = Depends(verify_session),
) -> UploadOverviewResponse | JSONResponse:
"""B03 재접속 시 업로드 현황(정본) — 완료 파일 목록 + 중단 세션 + 완료 여부.
localStorage 기반 표시는 캐시를 지우거나 다른 PC로 가면 사라진다(2026-08-04 사용자
보고). 화면은 진입 시 이 응답을 정본으로 삼고 localStorage는 보조로만 쓴다.
"""
pool = get_db_pool()
try:
async with pool.acquire() as connection:
files = await list_project_input_files(connection, project_id)
sessions = await list_incomplete_upload_sessions(connection, project_id)
file_types, point_cloud_id = await get_project_input_readiness(connection, project_id)
async with connection.cursor(aiomysql.DictCursor) as cursor:
state = await get_workflow_state(cursor, str(project_id))
stages = (state or {}).get("stages") or []
analysis_complete = any(
int(stage.get("stage_no", -1)) == 1 and str(stage.get("state")) == "COMPLETE"
for stage in stages
)
return UploadOverviewResponse(
files=[
UploadOverviewFile(
input_file_id=int(row["id"]),
file_type=str(row["file_type"]),
original_filename=str(row["original_filename"]),
file_size_mb=float(row["file_size_mb"] or 0.0),
status=str(row["status"]),
uploaded_at=str(row["upload_at"]) if row.get("upload_at") else None,
)
for row in files
],
pending_sessions=[
UploadOverviewSession(
upload_session_id=str(row["id"]),
original_filename=str(row["original_filename"]),
file_size_bytes=int(row["file_size_bytes"]),
total_chunks=int(row["total_chunks"]),
completed_chunks=int(row["completed_chunks"]),
progress_percent=round(
100.0 * int(row["completed_chunks"]) / max(1, int(row["total_chunks"])), 1
),
updated_at=str(row["updated_at"]) if row.get("updated_at") else None,
)
for row in sessions
],
required_complete=_REQUIRED_FILE_TYPES <= file_types and point_cloud_id is not None,
analysis_complete=analysis_complete,
)
except Exception:
logger.exception("B03 업로드 현황 조회 실패: project_id=%s", project_id)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "업로드 현황 조회 중 오류가 발생했습니다."},
)
@router.get("/{project_id}/workflow-state")
async def get_project_workflow_state(project_id: str):
pool = get_db_pool()
async with pool.acquire() as connection, connection.cursor(aiomysql.DictCursor) as cursor:
state = await get_workflow_state(cursor, project_id)
return {"status": "success", "workflow_state": state}