계획서 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>
404 lines
17 KiB
Python
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": "업로드 상태 조회 중 오류가 발생했습니다."},
|
|
)
|