diff --git a/B03_FileInput/B03_FileInput_Router.py b/B03_FileInput/B03_FileInput_Router.py index 17dd68b5..5728b6e9 100644 --- a/B03_FileInput/B03_FileInput_Router.py +++ b/B03_FileInput/B03_FileInput_Router.py @@ -53,11 +53,14 @@ from B03_FileInput.B03_FileInput_Schema import ( 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_json import atomic_write_json +from common_util.common_util_project_reset import purge_project_outputs from common_util.common_util_storage import resolve_stored_project_path from common_util.common_util_workflow import load_project_workflow from common_util.common_util_workflow_state import ( complete_stage, get_workflow_state, + is_analysis_running, + reset_stages_after_input_change, ) from config.config_db import get_db_pool from config.config_system import ( @@ -65,6 +68,8 @@ from config.config_system import ( UPLOAD_MAX_FILES, ) +_ANALYSIS_RUNNING_MESSAGE = "이 프로젝트는 지금 분석 중입니다. 끝난 뒤에 새 자료를 올려 주세요." + logger = logging.getLogger(__name__) router = APIRouter(prefix="/api/projects", tags=["B03 File Input"]) @@ -106,7 +111,14 @@ async def _complete_file_input_if_ready( _require_complete_file_set(file_types) if point_cloud_input_id is None: raise ValueError("LAS 또는 LAZ 입력 파일을 찾을 수 없습니다.") - async with connection.cursor() as cursor: + # 자료가 갈렸으니 옛 계산 결과(파일 + DB)를 지우고 진행 표시도 되돌린다. 남겨 두면 + # 아직 다시 만들어지지 않은 뒷단계 화면이 옛 결과를 새 자료 것인 양 보여준다. + stored_path = await get_project_storage_relative_path(connection, project_id) + await purge_project_outputs( + connection, str(project_id), Path(resolve_stored_project_path(stored_path)) + ) + 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) return point_cloud_input_id @@ -381,6 +393,14 @@ async def create_project_upload_session( 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}, + ) await create_upload_session( connection, session_id=session_id, diff --git a/B03_FileInput/B03_FileInput_Router_Temp.py b/B03_FileInput/B03_FileInput_Router_Temp.py index 5b1f7930..4078c170 100644 --- a/B03_FileInput/B03_FileInput_Router_Temp.py +++ b/B03_FileInput/B03_FileInput_Router_Temp.py @@ -12,6 +12,7 @@ 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 @@ -72,11 +73,16 @@ from B03_FileInput.B03_FileInput_Schema_Temp import ( ) 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 +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, @@ -84,6 +90,8 @@ from config.config_system import ( 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"]) @@ -611,6 +619,12 @@ async def attach_temp_batch( 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( @@ -634,6 +648,17 @@ async def attach_temp_batch( 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"]) @@ -668,7 +693,8 @@ async def attach_temp_batch( ) if str(item["file_type"]).lower() in _POINT_CLOUD_FILE_TYPES: point_cloud_input_id = input_file_id - async with connection.cursor() as cursor: + 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) diff --git a/B04_PreProcess/B04_PreProcess_Engine_Structurize.py b/B04_PreProcess/B04_PreProcess_Engine_Structurize.py index d4927f7b..f9e38b4b 100644 --- a/B04_PreProcess/B04_PreProcess_Engine_Structurize.py +++ b/B04_PreProcess/B04_PreProcess_Engine_Structurize.py @@ -8,6 +8,7 @@ from pathlib import Path import laspy import numpy as np +from common_util.common_util_json import replace_with_retry from config.config_system import SURFACE_DEFAULT_RGB_VALUE, SURFACE_LAS_CHUNK_SIZE @@ -101,7 +102,8 @@ def structurize_las( ) temporary.flush() os.fsync(temporary.fileno()) - os.replace(temporary_path, target) + # 이 교체가 윈도우 공유 위반으로 죽으면 WF1 분석 전체가 실패한다 — 재시도한다. + replace_with_retry(temporary_path, target) temporary_path = None finally: if temporary_path is not None: diff --git a/common_util/common_util_json.py b/common_util/common_util_json.py index 123c98de..912e59dc 100644 --- a/common_util/common_util_json.py +++ b/common_util/common_util_json.py @@ -3,9 +3,30 @@ import json import os import tempfile +import time from pathlib import Path from typing import Any +# 윈도우에서 `os.replace`는 대상 파일을 **누가 열고만 있어도** 거부당한다(WinError 5). +# 진행률 파일처럼 서버가 쓰는 동안 화면이 계속 읽는 파일에서 실제로 부딪혔고, 전처리 +# 산출물(structured.npz) 교체에서는 분석이 통째로 죽기도 했다(2026-08-08). +# 상대가 파일을 놓는 데 걸리는 시간은 밀리초 수준이라 짧게 여러 번 다시 시도한다. +_REPLACE_RETRY_COUNT = 12 +_REPLACE_RETRY_DELAY_SECONDS = 0.05 + + +def replace_with_retry(source: str | Path, target: str | Path) -> None: + """`os.replace` — 다른 쪽이 잠깐 잡고 있으면 잠시 기다렸다 다시 시도한다.""" + last_error: OSError | None = None + for attempt in range(_REPLACE_RETRY_COUNT): + try: + os.replace(source, target) + return + except PermissionError as error: # 윈도우 공유 위반 + last_error = error + time.sleep(_REPLACE_RETRY_DELAY_SECONDS * (attempt + 1)) + raise last_error if last_error else OSError("파일 교체에 실패했습니다.") + def atomic_write_json(path: str | Path, value: Any) -> None: """같은 디렉터리의 임시 파일을 교체하여 JSON을 원자적으로 저장한다.""" @@ -28,7 +49,7 @@ def atomic_write_json(path: str | Path, value: Any) -> None: os.fsync(temporary.fileno()) temporary_path = Path(temporary.name) - os.replace(temporary_path, target) + replace_with_retry(temporary_path, target) temporary_path = None finally: if temporary_path is not None: diff --git a/common_util/common_util_project_reset.py b/common_util/common_util_project_reset.py new file mode 100644 index 00000000..264b02c9 --- /dev/null +++ b/common_util/common_util_project_reset.py @@ -0,0 +1,96 @@ +"""입력 자료가 갈릴 때 옛 계산 결과를 지우는 유틸. + +새 자료를 올리면 B04~B06 자동 계산이 처음부터 다시 돌아 결과를 덮어쓴다. 그런데 옛 자료로 +만든 산출물이 남아 있으면, 아직 다시 만들어지지 않은 뒷단계(B07 이후) 화면이 **옛 결과를 +그대로 보여준다** — 사용자는 새 자료 기준인 줄 알고 검토하게 된다. + +그래서 자료가 갈리는 순간 계산 결과를 통째로 지운다. 화면 쪽 규칙은 단순해진다: +**불러올 자료가 없으면 대시보드로 돌려보낸다**(2026-08-08 사용자 지시). + +지우지 않는 것: `B03_FileInput/` — 업로드 원본과 그 메타데이터다. 같은 파일을 다시 올렸는지 +가리는 중복 검사가 이걸 본다. +""" + +import logging +import shutil +from pathlib import Path + +import aiomysql + +logger = logging.getLogger(__name__) + +# 계산 결과가 쌓이는 단계 폴더. B03(원본 입력)은 일부러 뺐다. +OUTPUT_STAGE_DIRS = ( + "B04_PreProcess", + "B05_Profile", + "B06_Section", + "B07_Quantity", + "B08_DesignDetail", + "B09_Estimation", +) + +# 프로젝트에 직접 매달린 산출물 테이블. 지우는 순서는 자식 → 부모. +_PROJECT_OUTPUT_TABLES = ( + "cross_sections", + "longitudinal_sections", + "structures", + "quantity_items", + "outputs", + "surface_models", + "processed_point_cloud", +) + + +def purge_output_folders(project_root: Path) -> list[str]: + """단계별 산출물 폴더를 비운다. 지운 폴더 이름을 돌려준다.""" + removed: list[str] = [] + for name in OUTPUT_STAGE_DIRS: + target = project_root / name + if not target.exists(): + continue + shutil.rmtree(target, ignore_errors=True) + target.mkdir(parents=True, exist_ok=True) + removed.append(name) + return removed + + +async def purge_output_records(connection: aiomysql.Connection, project_id: str) -> dict[str, int]: + """산출물 DB 레코드를 지운다. 파일만 지우면 화면은 "자료가 있다"고 착각한다. + + 호출부가 트랜잭션을 관리한다(여기서 커밋하지 않는다). + """ + deleted: dict[str, int] = {} + async with connection.cursor() as cursor: + # 노선에 매달린 자식부터 — 외래키 없이도 고아 행이 남지 않게 한다. + await cursor.execute( + "DELETE rp FROM route_points rp JOIN routes r ON r.id = rp.route_id " + "WHERE r.project_id = %s", + (project_id,), + ) + deleted["route_points"] = cursor.rowcount + await cursor.execute( + "DELETE rs FROM route_statistics rs JOIN routes r ON r.id = rs.route_id " + "WHERE r.project_id = %s", + (project_id,), + ) + deleted["route_statistics"] = cursor.rowcount + for table in _PROJECT_OUTPUT_TABLES: + await cursor.execute(f"DELETE FROM {table} WHERE project_id = %s", (project_id,)) + deleted[table] = cursor.rowcount + await cursor.execute("DELETE FROM routes WHERE project_id = %s", (project_id,)) + deleted["routes"] = cursor.rowcount + return deleted + + +async def purge_project_outputs( + connection: aiomysql.Connection, project_id: str, project_root: Path +) -> None: + """입력 자료 교체 시 옛 계산 결과(파일 + DB)를 함께 지운다.""" + removed = purge_output_folders(project_root) + deleted = await purge_output_records(connection, project_id) + logger.info( + "입력 자료 교체 — 옛 산출물 정리: project_id=%s 폴더=%s 레코드=%s", + project_id, + ",".join(removed) or "없음", + {key: value for key, value in deleted.items() if value}, + ) diff --git a/common_util/common_util_workflow_state.py b/common_util/common_util_workflow_state.py index 864eac59..a06e9280 100644 --- a/common_util/common_util_workflow_state.py +++ b/common_util/common_util_workflow_state.py @@ -219,4 +219,51 @@ async def get_workflow_state(cursor: aiomysql.DictCursor, project_id: str) -> Di else: break - return {"project_id": project_id, "current_stage": current_stage, "stages": stages_list} + return { + "project_id": project_id, + "current_stage": current_stage, + "stages": stages_list, + } + + +async def is_analysis_running(cursor: aiomysql.DictCursor, project_id: str) -> bool: + """전처리(1단계)가 도는 중인지. 도는 동안 새 자료를 받으면 분석 2개가 같은 산출물 + 경로를 놓고 부딪혀 뒤에 시작한 쪽이 죽는다(2026-08-08 E2E 점검에서 실제 발생). + + 화면은 업로드 버튼을 잠가 이 상황을 막지만, 새로고침·다른 탭·보관함 연결로 우회할 수 + 있어 서버에서도 막는다. + """ + await cursor.execute( + """ + SELECT state + FROM project_workflow_stages + WHERE project_id = %s AND stage_no = 1 + """, + (project_id,), + ) + row = await cursor.fetchone() + if not row: + return False + state = row["state"] if isinstance(row, dict) else row[0] + return str(state) == "IN_PROGRESS" + + +async def reset_stages_after_input_change(cursor: aiomysql.DictCursor, project_id: str) -> None: + """입력 자료가 갈렸을 때 1단계(전처리) 이후를 처음 상태로 되돌린다. + + 새 자료로 다시 계산해 덮어쓰므로, 옛 자료 기준의 진행 표시가 남아 있으면 사용자가 이미 + 끝난 단계로 착각한다. 파일 입력(0단계)은 방금 끝났으니 건드리지 않는다. + """ + await cursor.execute( + """ + UPDATE project_workflow_stages + SET state = 'NOT_STARTED', + progress_percent = 0, + params = NULL, + message = NULL, + started_at = NULL, + completed_at = NULL + WHERE project_id = %s AND stage_no >= 1 + """, + (project_id,), + )