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>
295 lines
12 KiB
Python
295 lines
12 KiB
Python
"""임시 보관함 청크 업로드 엔드포인트 — 세션 생성·조각 전송·진행 조회·마무리.
|
|
|
|
라우터가 700줄을 넘어 떼어냈다(2026-09-04). 본체가 `include_router` 로 붙이므로
|
|
여기 라우터에는 prefix 를 두지 않는다 — 두면 `/api/temp-uploads` 가 두 번 붙는다.
|
|
"""
|
|
|
|
import asyncio
|
|
import logging
|
|
from pathlib import Path
|
|
from typing import Any
|
|
from uuid import uuid4
|
|
|
|
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 (
|
|
list_completed_chunk_indexes,
|
|
mark_upload_session_completed,
|
|
mark_upload_session_failed,
|
|
upsert_upload_chunk,
|
|
)
|
|
from B03_FileInput.B03_FileInput_Repository_Temp import (
|
|
create_temp_upload_session,
|
|
get_temp_batch,
|
|
get_temp_upload_session,
|
|
upsert_temp_batch_file,
|
|
)
|
|
from B03_FileInput.B03_FileInput_Router_Errors import lookup_error_response
|
|
from B03_FileInput.B03_FileInput_Router_Temp_Support import (
|
|
_batch_root,
|
|
_refresh_batch_status,
|
|
_total_chunks,
|
|
)
|
|
from B03_FileInput.B03_FileInput_Schema import (
|
|
ChunkSessionCreateRequest,
|
|
ChunkSessionCreateResponse,
|
|
ChunkUploadResponse,
|
|
FileUploadDescriptor,
|
|
UploadFinalizeRequest,
|
|
UploadStatusResponse,
|
|
)
|
|
from B03_FileInput.B03_FileInput_Schema_Temp import (
|
|
TempFileUploadResponse,
|
|
TempFileUploadResult,
|
|
)
|
|
from common_util.common_util_auth import verify_session
|
|
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 Temp Upload"])
|
|
|
|
|
|
@router.post("/{batch_id}/upload-sessions", response_model=ChunkSessionCreateResponse)
|
|
async def create_batch_upload_session(
|
|
batch_id: str,
|
|
payload: ChunkSessionCreateRequest,
|
|
session: dict[str, Any] = Depends(verify_session),
|
|
) -> ChunkSessionCreateResponse | JSONResponse:
|
|
"""대용량 파일(LAS/LAZ) 청크 세션을 만든다 — 프로젝트 업로드와 같은 규칙."""
|
|
user_id = int(session["user_id"])
|
|
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 _batch_root(connection, batch_id=batch_id, user_id=user_id)
|
|
await create_temp_upload_session(
|
|
connection,
|
|
session_id=session_id,
|
|
batch_id=batch_id,
|
|
original_filename=payload.original_filename,
|
|
file_size_bytes=payload.size_bytes,
|
|
chunk_size_bytes=chunk_size_bytes,
|
|
total_chunks=total_chunks,
|
|
)
|
|
await connection.commit()
|
|
return ChunkSessionCreateResponse(
|
|
project_id=batch_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 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": "업로드 세션 생성 중 오류가 발생했습니다."},
|
|
)
|
|
|
|
|
|
@router.post("/{batch_id}/chunks", response_model=ChunkUploadResponse)
|
|
async def upload_batch_chunk(
|
|
batch_id: str,
|
|
session_id: str = Form(...),
|
|
chunk_index: int = Form(...),
|
|
chunk_data: UploadFile = File(...),
|
|
session: dict[str, Any] = Depends(verify_session),
|
|
) -> ChunkUploadResponse | JSONResponse:
|
|
"""청크 한 조각을 보관함 묶음 폴더에 저장한다."""
|
|
user_id = int(session["user_id"])
|
|
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)
|
|
upload_session = await get_temp_upload_session(
|
|
connection, batch_id=batch_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": "청크 인덱스가 범위를 벗어났습니다."},
|
|
)
|
|
session_dir = resolve_chunk_session_dir(batch_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"]),
|
|
)
|
|
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=chunk_path.relative_to(batch_root).as_posix(),
|
|
)
|
|
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 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:
|
|
await chunk_data.close()
|
|
|
|
|
|
@router.get("/{batch_id}/upload-status/{session_id}", response_model=UploadStatusResponse)
|
|
async def get_batch_upload_status(
|
|
batch_id: str,
|
|
session_id: str,
|
|
session: dict[str, Any] = Depends(verify_session),
|
|
) -> UploadStatusResponse | 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)
|
|
upload_session = await get_temp_upload_session(
|
|
connection, batch_id=batch_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 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.post("/{batch_id}/finalize", response_model=TempFileUploadResponse)
|
|
async def finalize_batch_upload(
|
|
batch_id: str,
|
|
payload: UploadFinalizeRequest,
|
|
session: dict[str, Any] = Depends(verify_session),
|
|
) -> TempFileUploadResponse | JSONResponse:
|
|
"""청크를 병합해 보관함에 저장하고, 필수 파일이 다 차면 완료로 올린다."""
|
|
user_id = int(session["user_id"])
|
|
pool = get_db_pool()
|
|
final_path: Path | None = None
|
|
try:
|
|
async with pool.acquire() as connection:
|
|
batch_root = await _batch_root(connection, batch_id=batch_id, user_id=user_id)
|
|
upload_session = await get_temp_upload_session(
|
|
connection, batch_id=batch_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
|
|
)
|
|
if completed_indexes != list(range(payload.total_chunks)):
|
|
return JSONResponse(
|
|
status_code=400,
|
|
content={"status": "error", "message": "아직 업로드되지 않은 청크가 있습니다."},
|
|
)
|
|
|
|
descriptor = FileUploadDescriptor(
|
|
original_filename=str(upload_session["original_filename"]),
|
|
size_bytes=int(upload_session["file_size_bytes"]),
|
|
)
|
|
final_path = merge_upload_chunks(
|
|
batch_root, descriptor, payload.session_id, payload.total_chunks
|
|
)
|
|
metadata = await asyncio.to_thread(analyze_input_metadata, final_path)
|
|
relative_path = final_path.relative_to(batch_root).as_posix()
|
|
file_type = final_path.suffix.lower().lstrip(".")
|
|
crs_epsg = metadata.get("epsg")
|
|
|
|
await connection.begin()
|
|
try:
|
|
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=int(upload_session["file_size_bytes"]),
|
|
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)
|
|
required_complete = await _refresh_batch_status(connection, batch_id=batch_id)
|
|
await connection.commit()
|
|
except Exception:
|
|
await connection.rollback()
|
|
await mark_upload_session_failed(connection, session_id=payload.session_id)
|
|
raise
|
|
|
|
remove_chunk_session(batch_root, payload.session_id)
|
|
return TempFileUploadResponse(
|
|
batch_id=batch_id,
|
|
files=[
|
|
TempFileUploadResult(
|
|
batch_id=batch_id,
|
|
file_type=file_type,
|
|
original_filename=descriptor.original_filename,
|
|
relative_path=relative_path,
|
|
size_bytes=int(upload_session["file_size_bytes"]),
|
|
metadata=metadata,
|
|
)
|
|
],
|
|
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:
|
|
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("임시 보관함 병합 실패: batch_id=%s", batch_id)
|
|
return JSONResponse(
|
|
status_code=500,
|
|
content={"status": "error", "message": "업로드 최종 처리 중 오류가 발생했습니다."},
|
|
)
|