749줄 한 파일을 셋으로 나눔 (동작 불변, 순수 이동). - `B03_FileInput_Router_Temp.py` 478줄 — 묶음 생성·목록·삭제·파일 업로드·프로젝트 연결 - `B03_FileInput_Router_Temp_Chunks.py` 294줄 — 세션·조각·진행·마무리 (prefix 없는 APIRouter 를 본체가 `include_router` 로 붙여 `/api/temp-uploads/...` 경로 불변) - `B03_FileInput_Router_Temp_Support.py` 46줄 — 조각 수·묶음 폴더·상태 갱신 (순환 방지) 검증: 라우트 10개(temp 9 + attach 1) 경로·메서드 동일, 공용 브라우저에서 실제 청크 흐름 create 200 → upload-sessions 200 → chunks 200 → finalize 200 → delete 200, ruff check 통과, tmp/tests 378 passed Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
479 lines
21 KiB
Python
479 lines
21 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
|
|
from B03_FileInput.B03_FileInput_Router_Temp_Chunks import router as chunk_router
|
|
from B03_FileInput.B03_FileInput_Router_Temp_Support import (
|
|
_POINT_CLOUD_FILE_TYPES,
|
|
_batch_root,
|
|
_iso,
|
|
_refresh_batch_status,
|
|
)
|
|
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 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,
|
|
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)를 먼저 지운다 — 남겨 두면 아직 다시
|
|
# 만들어지지 않은 뒷단계 화면이 옛 결과를 새 자료 것인 양 보여준다.
|
|
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": "보관함 연결 중 오류가 발생했습니다."},
|
|
)
|