Files
Aislo/B03_FileInput/B03_FileInput_Router_Temp.py
T
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

513 lines
22 KiB
Python

"""임시 보관함 라우터 — 프로젝트 생성 전에 계정에 묶어 자료를 올려 둔다.
라이다 원본은 업로드에 오래 걸려서 프로젝트 정보가 확정되기 전에 미리 올릴 수 있어야
한다(2026-08-08 사용자 지시). 저장 구조를 프로젝트 저장소와 똑같이 맞춰 두어 청크
업로드·병합 엔진을 그대로 재사용하고, 프로젝트로 옮길 때도 같은 상대 경로로 붙인다.
"""
import asyncio
import logging
import shutil
from pathlib import Path
from typing import Any
from uuid import UUID, uuid4
import aiomysql
from fastapi import APIRouter, Depends, File, UploadFile
from fastapi.responses import JSONResponse
from B03_FileInput.B03_FileInput_Engine import (
resolve_upload_destination,
save_upload_stream,
)
from B03_FileInput.B03_FileInput_Engine_Analyze import analyze_input_metadata
from B03_FileInput.B03_FileInput_Repository import (
create_input_file,
get_project_storage_relative_path,
supersede_previous_input_files,
)
from B03_FileInput.B03_FileInput_Repository_Temp import (
create_temp_batch,
delete_temp_batch,
delete_temp_batch_file,
get_temp_batch,
get_temp_batch_file,
is_batch_required_complete,
list_temp_batch_files,
list_temp_batch_sessions,
list_temp_batches,
mark_temp_batch_linked,
upsert_temp_batch_file,
)
from B03_FileInput.B03_FileInput_Router_Errors import lookup_error_response
# 분리 전 이 파일에 있던 이름은 그대로 다시 내보낸다(2026-09-04).
from B03_FileInput.B03_FileInput_Router_Temp_Chunks import (
create_batch_upload_session as create_batch_upload_session,
)
from B03_FileInput.B03_FileInput_Router_Temp_Chunks import (
finalize_batch_upload as finalize_batch_upload,
)
from B03_FileInput.B03_FileInput_Router_Temp_Chunks import (
get_batch_upload_status as get_batch_upload_status,
)
from B03_FileInput.B03_FileInput_Router_Temp_Chunks import router as chunk_router
from B03_FileInput.B03_FileInput_Router_Temp_Chunks import (
upload_batch_chunk as upload_batch_chunk,
)
from B03_FileInput.B03_FileInput_Router_Temp_Support import (
_POINT_CLOUD_FILE_TYPES as _POINT_CLOUD_FILE_TYPES,
)
from B03_FileInput.B03_FileInput_Router_Temp_Support import (
_batch_root as _batch_root,
)
from B03_FileInput.B03_FileInput_Router_Temp_Support import (
_iso as _iso,
)
from B03_FileInput.B03_FileInput_Router_Temp_Support import (
_refresh_batch_status as _refresh_batch_status,
)
from B03_FileInput.B03_FileInput_Router_Temp_Support import (
_total_chunks as _total_chunks,
)
from B03_FileInput.B03_FileInput_Schema import (
FileUploadDescriptor,
)
from B03_FileInput.B03_FileInput_Schema_Temp import (
TempBatchAttachResponse,
TempBatchCreateRequest,
TempBatchCreateResponse,
TempBatchFile,
TempBatchItem,
TempBatchListResponse,
TempBatchPendingSession,
TempFileUploadResponse,
TempFileUploadResult,
)
from B03_FileInput.B03_FileInput_Service_WF1 import trigger_wf1_analysis_and_email
from common_util.common_util_auth import verify_session
from B03_FileInput.B03_FileInput_Router_Helpers import (
OutputsWouldBeDiscarded,
_confirm_replace_response,
describe_existing_outputs,
)
from common_util.common_util_project_reset import purge_project_outputs
from common_util.common_util_storage import (
resolve_stored_project_path,
resolve_temp_batch_path,
)
from common_util.common_util_workflow_state import (
complete_stage,
is_analysis_running,
reset_stages_after_input_change,
)
from config.config_db import get_db_pool
from config.config_system import (
TEMP_UPLOAD_RETENTION_DAYS,
UPLOAD_MAX_FILES,
)
_ANALYSIS_RUNNING_MESSAGE = "이 프로젝트는 지금 분석 중입니다. 끝난 뒤에 다시 시도해 주세요."
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/temp-uploads", tags=["B03 Temp Upload"])
attach_router = APIRouter(prefix="/api/projects", tags=["B03 Temp Upload"])
# 청크 업로드 엔드포인트는 파일이 700줄을 넘어 떼어냈다(2026-09-04) — 경로는 그대로다.
router.include_router(chunk_router)
@router.post("", response_model=TempBatchCreateResponse)
async def create_batch(
payload: TempBatchCreateRequest,
session: dict[str, Any] = Depends(verify_session),
) -> TempBatchCreateResponse | JSONResponse:
"""보관함 묶음(파일 한 세트)을 만든다."""
batch_id = str(uuid4())
user_id = int(session["user_id"])
pool = get_db_pool()
try:
async with pool.acquire() as connection:
await create_temp_batch(
connection,
batch_id=batch_id,
user_id=user_id,
name=payload.name,
memo=payload.memo,
)
await connection.commit()
resolve_temp_batch_path(user_id, batch_id)
return TempBatchCreateResponse(batch_id=batch_id, name=payload.name)
except Exception:
logger.exception("임시 보관함 묶음 생성 실패: user_id=%s", user_id)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "보관함 생성 중 오류가 발생했습니다."},
)
@router.get("", response_model=TempBatchListResponse)
async def list_batches(
session: dict[str, Any] = Depends(verify_session),
) -> TempBatchListResponse | JSONResponse:
"""내 보관함 목록 — 파일 목록과 진행 중 세션 진행률을 함께 준다."""
user_id = int(session["user_id"])
pool = get_db_pool()
try:
async with pool.acquire() as connection:
batches = await list_temp_batches(connection, user_id=user_id)
ids = [str(batch["id"]) for batch in batches]
files_by_batch = await list_temp_batch_files(connection, batch_ids=ids)
sessions_by_batch = await list_temp_batch_sessions(connection, batch_ids=ids)
items: list[TempBatchItem] = []
for batch in batches:
batch_id = str(batch["id"])
files = files_by_batch.get(batch_id, [])
sessions = sessions_by_batch.get(batch_id, [])
file_types = {str(item["file_type"]).lower() for item in files}
items.append(
TempBatchItem(
batch_id=batch_id,
name=str(batch["name"]),
memo=batch.get("memo"),
status=str(batch["status"]),
files=[
TempBatchFile(
file_type=str(item["file_type"]),
original_filename=str(item["original_filename"]),
file_size_bytes=int(item["file_size_bytes"]),
crs_epsg=item.get("crs_epsg"),
)
for item in files
],
pending_sessions=[
TempBatchPendingSession(
upload_session_id=str(item["id"]),
original_filename=str(item["original_filename"]),
file_size_bytes=int(item["file_size_bytes"]),
total_chunks=int(item["total_chunks"]),
completed_chunks=int(item["completed_chunks"]),
progress_percent=round(
int(item["completed_chunks"])
/ max(1, int(item["total_chunks"]))
* 100,
1,
),
)
for item in sessions
],
total_size_bytes=sum(int(item["file_size_bytes"]) for item in files),
required_complete=is_batch_required_complete(file_types),
completed_at=_iso(batch.get("completed_at")),
expires_at=_iso(batch.get("expires_at")),
linked_project_id=(
str(batch["linked_project_id"]) if batch.get("linked_project_id") else None
),
created_at=_iso(batch.get("created_at")),
)
)
return TempBatchListResponse(batches=items, retention_days=TEMP_UPLOAD_RETENTION_DAYS)
except Exception:
logger.exception("임시 보관함 목록 조회 실패: user_id=%s", user_id)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "보관함 조회 중 오류가 발생했습니다."},
)
@router.delete("/{batch_id}")
async def remove_batch(
batch_id: str,
session: dict[str, Any] = Depends(verify_session),
) -> JSONResponse:
"""보관함 묶음과 저장 파일을 지운다."""
user_id = int(session["user_id"])
pool = get_db_pool()
try:
async with pool.acquire() as connection:
await get_temp_batch(connection, batch_id=batch_id, user_id=user_id)
await delete_temp_batch(connection, batch_id=batch_id, user_id=user_id)
await connection.commit()
batch_root = Path(resolve_temp_batch_path(user_id, batch_id, create=False))
shutil.rmtree(batch_root, ignore_errors=True)
return JSONResponse(content={"status": "success", "batch_id": batch_id})
except LookupError as exc:
return lookup_error_response(exc, logger, context="B03 보관함", batch_id=batch_id)
except Exception:
logger.exception("임시 보관함 삭제 실패: batch_id=%s", batch_id)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "보관함 삭제 중 오류가 발생했습니다."},
)
@router.delete("/{batch_id}/files/{file_type}")
async def remove_batch_file(
batch_id: str,
file_type: str,
session: dict[str, Any] = Depends(verify_session),
) -> JSONResponse:
"""묶음 안 파일 1건을 지운다 — DB 행과 임시 저장소 실제 파일을 함께 지운다."""
user_id = int(session["user_id"])
normalized = file_type.lower().lstrip(".")
pool = get_db_pool()
try:
async with pool.acquire() as connection:
batch = await get_temp_batch(connection, batch_id=batch_id, user_id=user_id)
if str(batch["status"]) == "linked":
return JSONResponse(
status_code=400,
content={
"status": "error",
"message": "프로젝트에 연결된 보관함은 수정할 수 없습니다.",
},
)
item = await get_temp_batch_file(connection, batch_id=batch_id, file_type=normalized)
await delete_temp_batch_file(connection, batch_id=batch_id, file_type=normalized)
required_complete = await _refresh_batch_status(connection, batch_id=batch_id)
await connection.commit()
batch_root = Path(resolve_temp_batch_path(user_id, batch_id, create=False))
stored = batch_root / str(item["relative_path"])
stored.unlink(missing_ok=True)
return JSONResponse(
content={
"status": "success",
"batch_id": batch_id,
"file_type": normalized,
"required_complete": required_complete,
}
)
except LookupError as exc:
return lookup_error_response(exc, logger, context="B03 보관함", batch_id=batch_id)
except OSError as exc:
return JSONResponse(status_code=400, content={"status": "error", "message": str(exc)})
except Exception:
logger.exception("보관함 파일 삭제 실패: batch_id=%s type=%s", batch_id, file_type)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "파일 삭제 중 오류가 발생했습니다."},
)
@router.post("/{batch_id}/files", response_model=TempFileUploadResponse)
async def upload_batch_files(
batch_id: str,
files: list[UploadFile] = File(...),
session: dict[str, Any] = Depends(verify_session),
) -> TempFileUploadResponse | JSONResponse:
"""작은 파일(csv·prj·tfw·tif)을 보관함에 바로 저장한다."""
user_id = int(session["user_id"])
if not files or len(files) > UPLOAD_MAX_FILES:
message = f"파일은 1~{UPLOAD_MAX_FILES}개까지 가능합니다."
return JSONResponse(status_code=400, content={"status": "error", "message": message})
pool = get_db_pool()
try:
async with pool.acquire() as connection:
batch_root = await _batch_root(connection, batch_id=batch_id, user_id=user_id)
results: list[TempFileUploadResult] = []
for upload in files:
descriptor = FileUploadDescriptor(
original_filename=upload.filename or "",
size_bytes=max(1, upload.size or 1),
)
destination = resolve_upload_destination(batch_root, descriptor)
written_bytes = await save_upload_stream(upload, destination)
metadata = await asyncio.to_thread(analyze_input_metadata, destination)
relative_path = destination.relative_to(batch_root).as_posix()
file_type = destination.suffix.lower().lstrip(".")
crs_epsg = metadata.get("epsg")
await upsert_temp_batch_file(
connection,
batch_id=batch_id,
file_type=file_type,
original_filename=descriptor.original_filename,
relative_path=relative_path,
file_size_bytes=written_bytes,
crs_epsg=int(crs_epsg) if crs_epsg is not None else None,
metadata=metadata,
)
results.append(
TempFileUploadResult(
batch_id=batch_id,
file_type=file_type,
original_filename=descriptor.original_filename,
relative_path=relative_path,
size_bytes=written_bytes,
metadata=metadata,
)
)
required_complete = await _refresh_batch_status(connection, batch_id=batch_id)
await connection.commit()
return TempFileUploadResponse(
batch_id=batch_id, files=results, required_complete=required_complete
)
except LookupError as exc:
return lookup_error_response(exc, logger, context="B03 보관함", batch_id=batch_id)
except (OSError, ValueError) as exc:
return JSONResponse(status_code=400, content={"status": "error", "message": str(exc)})
except Exception:
logger.exception("임시 보관함 파일 저장 실패: batch_id=%s", batch_id)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "파일 저장 중 오류가 발생했습니다."},
)
finally:
for upload in files:
await upload.close()
@attach_router.post(
"/{project_id}/temp-uploads/{batch_id}/attach", response_model=TempBatchAttachResponse
)
async def attach_temp_batch(
project_id: UUID,
batch_id: str,
confirm_replace: bool = False,
session: dict[str, Any] = Depends(verify_session),
) -> TempBatchAttachResponse | JSONResponse:
"""보관함 자료를 프로젝트 영구저장소로 옮기고 초기 분석을 시작한다.
파일 이동 → `input_files` 등록 → stage 0 완료 → WF1·자동 설계 체인까지, B03에서
직접 업로드했을 때와 같은 흐름을 탄다.
"""
from B03_FileInput.B03_FileInput_Router import _schedule_background_task
user_id = int(session["user_id"])
pool = get_db_pool()
try:
async with pool.acquire() as connection:
async with connection.cursor(aiomysql.DictCursor) as cursor:
if await is_analysis_running(cursor, str(project_id)):
return JSONResponse(
status_code=409,
content={"status": "error", "message": _ANALYSIS_RUNNING_MESSAGE},
)
batch = await get_temp_batch(connection, batch_id=batch_id, user_id=user_id)
if str(batch["status"]) == "linked":
return JSONResponse(
status_code=400,
content={"status": "error", "message": "이미 프로젝트에 연결된 보관함입니다."},
)
files = (await list_temp_batch_files(connection, batch_ids=[batch_id])).get(
batch_id, []
)
file_types = {str(item["file_type"]).lower() for item in files}
if not is_batch_required_complete(file_types):
return JSONResponse(
status_code=400,
content={
"status": "error",
"message": "필수 파일이 모두 갖춰진 보관함만 연결할 수 있습니다.",
},
)
stored_path = await get_project_storage_relative_path(connection, project_id)
batch_root = Path(resolve_temp_batch_path(user_id, batch_id, create=False))
project_root = Path(resolve_stored_project_path(stored_path))
# 자료가 갈리므로 옛 계산 결과(파일 + DB)를 먼저 지운다 — 남겨 두면 아직 다시
# 만들어지지 않은 뒷단계 화면이 옛 결과를 새 자료 것인 양 보여준다.
# ⚠ 되돌릴 수 없다 — 지울 것이 있으면 먼저 묻는다(2026-09-08).
if not confirm_replace:
targets = describe_existing_outputs(project_root)
if targets:
return _confirm_replace_response(OutputsWouldBeDiscarded(targets))
async with pool.acquire() as connection:
await connection.begin()
try:
await purge_project_outputs(connection, str(project_id), project_root)
await connection.commit()
except Exception:
await connection.rollback()
raise
moved: list[dict[str, Any]] = []
for item in files:
source = batch_root / str(item["relative_path"])
if not source.is_file():
name = item["original_filename"]
raise FileNotFoundError(f"보관함 파일을 찾을 수 없습니다: {name}")
destination = project_root / str(item["relative_path"])
destination.parent.mkdir(parents=True, exist_ok=True)
await asyncio.to_thread(shutil.move, str(source), str(destination))
moved.append({**item, "destination": destination})
point_cloud_input_id: int | None = None
async with pool.acquire() as connection:
await connection.begin()
try:
for item in moved:
metadata = item.get("metadata")
if isinstance(metadata, str):
import json as _json
metadata = _json.loads(metadata)
input_file_id = await create_input_file(
connection,
project_id=project_id,
file_type=str(item["file_type"]),
original_filename=str(item["original_filename"]),
relative_path=str(item["relative_path"]),
file_size_bytes=int(item["file_size_bytes"]),
upload_by=user_id,
crs_epsg=item.get("crs_epsg"),
metadata=metadata or {},
)
# 같은 이름의 옛 행은 내려 둔다 — 목록·분석이 최신 1건만 보게 한다.
await supersede_previous_input_files(
connection,
project_id,
str(item["original_filename"]),
input_file_id,
)
if str(item["file_type"]).lower() in _POINT_CLOUD_FILE_TYPES:
point_cloud_input_id = input_file_id
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)
await mark_temp_batch_linked(
connection, batch_id=batch_id, project_id=str(project_id)
)
await connection.commit()
except Exception:
await connection.rollback()
raise
shutil.rmtree(batch_root, ignore_errors=True)
analysis_started = point_cloud_input_id is not None
if analysis_started:
_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-preprocess-auto-{project_id}",
)
logger.info(
"보관함 연결 완료: project_id=%s batch_id=%s 파일=%d건 분석시작=%s",
project_id,
batch_id,
len(moved),
analysis_started,
)
return TempBatchAttachResponse(
project_id=str(project_id),
batch_id=batch_id,
moved_files=len(moved),
analysis_started=analysis_started,
)
except LookupError as exc:
return lookup_error_response(exc, logger, context="B03 보관함", project_id=project_id)
except (OSError, ValueError) as exc:
return JSONResponse(status_code=400, content={"status": "error", "message": str(exc)})
except Exception:
logger.exception("보관함 연결 실패: project_id=%s batch_id=%s", project_id, batch_id)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "보관함 연결 중 오류가 발생했습니다."},
)