Files
Aislo/B03_FileInput/B03_FileInput_Router_Chunks.py
T
eomsangdonandClaude Opus 5 018ebe7f40 feat(B03): 업로드 잠금에 이름을 붙임 — 지금 누가 올리는 중인지 화면에 뜸
계획서 0-8 의 남은 하나. 서버는 분석이 도는 동안 새 자료를 이미 막고 있었고
(`is_analysis_running`), 없던 것은 **막힌 까닭을 사람에게 보이는 것**이었음.
여럿이 한 프로젝트를 보면 「왜 안 올라가지」가 됨.

- 분석을 띄우는 **경로 넷 모두** 시작한 사람을 `params.started_by` 로 담음.
- `analysis_lock_owner` 가 그 id 로 이름을 붙여 냄 — 옛 자료는 id 가 없어
  **이름 없이 잠금만** 돎(막는 것이 먼저).
- `/upload-overview` 응답에 `analysis_lock` 을 실어 B03 화면이 띠로 보임.

화면 실측 — 「엄상돈 님이 올린 자료를 분석하는 중입니다」 띠가 경고색으로 뜸(응답을
잠깐 가로채 확인하고 곧바로 되돌림). 지금 도는 분석이 없어 실제 잠금 화면은 다음 업로드 때 볼 것.
시험 넷 추가(도는 중 아님·이름 붙음·옛 자료·경로 넷 대조), 전체 1211 통과·실패 0.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-09 22:43:28 +09:00

404 lines
17 KiB
Python

"""B03 청크 업로드 엔드포인트 — 세션 생성·조각 전송·마무리·진행 조회.
라우터 본체(`B03_FileInput_Router.py`)가 700줄을 넘어 떼어낸 조각이다(2026-09-04).
경로·태그는 본체와 같다 — 본체가 `include_router` 로 붙여 URL 이 그대로 유지된다.
"""
import asyncio
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_Engine import (
merge_upload_chunks,
remove_chunk_session,
resolve_chunk_session_dir,
save_upload_chunk,
)
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_storage_relative_path,
get_upload_session,
list_completed_chunk_indexes,
mark_upload_session_completed,
mark_upload_session_failed,
supersede_previous_input_files,
upsert_upload_chunk,
)
from B03_FileInput.B03_FileInput_Router_Errors import lookup_error_response
from B03_FileInput.B03_FileInput_Router_Helpers import (
_ANALYSIS_RUNNING_MESSAGE,
OutputsWouldBeDiscarded,
_already_uploaded,
_complete_file_input_if_ready,
_confirm_replace_response,
_schedule_background_task,
_total_chunks,
_write_stage_metadata,
)
from B03_FileInput.B03_FileInput_Schema import (
ChunkSessionCreateRequest,
ChunkSessionCreateResponse,
ChunkUploadResponse,
FileUploadDescriptor,
FileUploadResponse,
UploadedFileResult,
UploadFinalizeRequest,
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_storage import resolve_stored_project_path
from common_util.common_util_workflow_state import (
is_analysis_running,
)
from config.config_db import get_db_pool
from config.config_system import (
UPLOAD_CHUNK_SIZE_BYTES,
)
logger = logging.getLogger(__name__)
router = APIRouter(tags=["B03 File Input"])
@router.post("/{project_id}/upload-sessions", response_model=ChunkSessionCreateResponse)
async def create_project_upload_session(
project_id: UUID,
payload: ChunkSessionCreateRequest,
session: dict[str, Any] = Depends(verify_session),
) -> ChunkSessionCreateResponse | JSONResponse:
"""대용량 파일 청크 업로드 세션을 생성한다."""
# LAS 없는 설계를 켠 상태면 포인트클라우드는 받지 않는다 — 큰 LAS는 이 경로로
# 들어오므로 여기서 막지 않으면 `/files` 검사를 통째로 비켜 간다.
if payload.las_free and Path(payload.original_filename).suffix.lower() in {".las", ".laz"}:
return JSONResponse(
status_code=400,
content={
"status": "error",
"message": "LAS 없이 설계를 켠 상태에서는 LAS/LAZ 파일을 올릴 수 없습니다.",
},
)
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())
point_cloud_input_id: int | None = None
skipped: ChunkSessionCreateResponse | None = None
pool = get_db_pool()
try:
async with pool.acquire() as connection:
await get_project_storage_relative_path(connection, project_id)
# 분석이 도는 중이면 새 자료를 받지 않는다 — 받아 봐야 분석 2개가 같은 산출물
# 경로에서 부딪힌다. 화면도 버튼을 잠그지만 새로고침으로 우회할 수 있어 여기서 막는다.
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},
)
skipped = await _already_uploaded(connection, project_id, payload)
if skipped is not None:
if payload.complete_upload:
await connection.begin()
try:
point_cloud_input_id = await _complete_file_input_if_ready(
connection,
project_id,
payload.las_free,
payload.confirm_replace,
)
await connection.commit()
except Exception:
await connection.rollback()
raise
else:
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,
)
if skipped is not None:
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"]),
started_by=session.get("user_id"),
),
task_name=f"b04-preprocess-auto-{project_id}",
)
return skipped
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 OutputsWouldBeDiscarded as exc:
return _confirm_replace_response(exc)
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("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:
upload_session = await get_upload_session(
connection,
project_id=project_id,
session_id=session_id,
)
if chunk_index < 0 or chunk_index >= int(upload_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(upload_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(upload_session["total_chunks"]),
chunk_hash=chunk_hash,
)
except OutputsWouldBeDiscarded as exc:
return _confirm_replace_response(exc)
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(
"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:
upload_session = await get_upload_session(
connection,
project_id=project_id,
session_id=payload.session_id,
)
if int(upload_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(upload_session["original_filename"]),
size_bytes=int(upload_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)
# 다음에 같은 파일이 올라오면 전송을 건너뛸 수 있도록 지문을 함께 남긴다.
fingerprint = payload.fingerprint or None
if fingerprint:
metadata = {**metadata, "fingerprint": fingerprint}
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,
)
# 같은 이름의 옛 행은 내려 둔다 — 목록·분석이 최신 1건만 보게 한다.
await supersede_previous_input_files(
connection,
project_id,
descriptor.original_filename,
input_file_id,
)
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,
payload.las_free,
payload.confirm_replace,
)
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 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"]),
started_by=session.get("user_id"),
),
task_name=f"b04-preprocess-auto-{project_id}",
)
return FileUploadResponse(project_id=str(project_id), files=[result])
except OutputsWouldBeDiscarded as exc:
return _confirm_replace_response(exc)
except LookupError as exc:
return lookup_error_response(exc, logger, context="B03 업로드", project_id=project_id)
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:
upload_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(upload_session["status"]),
original_filename=str(upload_session["original_filename"]),
file_size_bytes=int(upload_session["file_size_bytes"]),
chunk_size_bytes=int(upload_session["chunk_size_bytes"]),
total_chunks=int(upload_session["total_chunks"]),
completed_chunks=len(completed_indexes),
completed_chunk_indexes=completed_indexes,
)
except OutputsWouldBeDiscarded as exc:
return _confirm_replace_response(exc)
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(
"B03 업로드 상태 조회 실패: project_id=%s session_id=%s",
project_id,
session_id,
)
return JSONResponse(
status_code=500,
content={"status": "error", "message": "업로드 상태 조회 중 오류가 발생했습니다."},
)