feat(B03): 자료 교체 시 옛 산출물 정리 + 분석 중 업로드 차단 (E2E 결함 3·6 서버측)
사용자 결정(2026-08-08): 새 자료를 올리면 재계산해 덮어쓰고, 화면은 불러올 자료가 없으면 대시보드로 보낸다. 그러려면 자료가 갈리는 순간 옛 계산 결과가 남아 있으면 안 된다. - common_util_project_reset 신설: 단계별 산출물 폴더(B04~B09)와 산출물 DB 레코드 (surface_models·routes·route_points·route_statistics·longitudinal_sections· cross_sections·structures·quantity_items·outputs·processed_point_cloud)를 함께 지운다. **B03_FileInput(업로드 원본)은 지우지 않는다** — 같은 파일인지 가리는 중복 검사가 쓴다. - 직접 업로드 완료·보관함 연결 양쪽에서 정리를 호출하고, 진행 단계도 1단계 이후를 NOT_STARTED로 되돌린다(reset_stages_after_input_change). - is_analysis_running(): 전처리가 도는 중이면 업로드 세션 생성과 보관함 연결을 409로 막는다. 화면 버튼 잠금은 새로고침·다른 탭으로 우회되므로 서버에도 문을 단다. - fail_stage 아래 있던 리비전(번호표) 안은 폐기 — 사용자가 "산출물이 없으면 대시보드" 방식으로 정리했다. 곁들여: os.replace 공유 위반 재시도(replace_with_retry) 진행률 파일은 서버가 쓰는 동안 화면이 계속 읽어 WinError 5가 났고, 전처리 structured.npz 교체에서는 같은 이유로 분석이 통째로 죽었다(결함 3의 사망 원인). 짧게 여러 번 다시 시도하도록 바꿨다. 검증(실서버 f45243b3, 계획노선 CSV 재업로드로 자료 교체): 전처리 결과 112개 -> 재분석분만, 노선 2->0, 횡단 20->0, DB 지표면 15->0 / 노선 1->0 / 횡단 19->0, 업로드 원본 6개는 그대로. 진행 단계가 초기로 돌아가고 재분석이 자동 시작됨. 교체 재시도는 읽는 쪽이 파일을 잡고 있는 상황을 만들어 성공 확인. ruff format·check 통과. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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},
|
||||
)
|
||||
@@ -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,),
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user