Files
Aislo/B03_FileInput/B03_FileInput_Service_WF1.py
T
eomsangdonandClaude Opus 5 e2f6a101b6 feat(B03): 재접속 업로드 현황(10) + B05·B06 자동 설계 체인(9)
10. 재접속 현황·완료 표시·재업로드 경고:
- GET /projects/{id}/upload-overview 신설 — 완료 파일 목록(input_files 정본),
  중단 청크 세션(진행률), 필수 파일·WF1 분석 완료 여부
- 진입 시 서버 정본으로 슬롯 카드 표시(serverUploaded), localStorage는 보조로 강등
- 파일 재선택 전에도 중단 세션 이어올리기 안내 배너
- 전체 완료 배지 + 완료 슬롯 재업로드 시 교체 확인 모달(승인 시에만 진행)
- 필수 슬롯 검증: 서버 업로드분 있으면 충족 — 단일 파일 교체 업로드 허용

9. 자동 설계 체인 연장 (B03_FileInput_Service_Chain.py):
- WF1 자동 확정 후 같은 백그라운드 태스크에서 ① 계획노선 CSV 기반 B05 기본 경로
  계산(solve) ② 경로 확정(stage 2) ③ B06 기본 횡단 설계 확정(stage 3)까지 진행
- 수동 이력 보호: 프로젝트에 경로가 이미 있으면 건너뜀
- 단계별 실패 격리: 실패 단계에서 멈추고 로그·workflow 상태로만 기록
- AUTO_DESIGN_CHAIN_ENABLED config 플래그(기본 True)

typecheck·ruff·B03 unittest(7건) 통과.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-04 18:09:15 +09:00

177 lines
7.4 KiB
Python

"""B03 업로드 이후 WF1 백그라운드 분석·자동 확정 서비스."""
import asyncio
import logging
from pathlib import Path
from typing import Any
from uuid import UUID
import aiomysql
from B03_FileInput.B03_FileInput_Email import (
send_analysis_completion_email,
send_analysis_error_email,
)
from B03_FileInput.B03_FileInput_Repository import get_project_storage_relative_path
from common_util.common_util_storage import resolve_stored_project_path
from common_util.common_util_surface_confirmation import surface_confirmation_defaults
from common_util.common_util_workflow_state import fail_stage, start_stage
from config.config_db import get_db_pool
from config.config_system import (
AUTO_DESIGN_CHAIN_ENABLED,
SEND_ANALYSIS_COMPLETION_EMAIL,
SURFACE_MODEL_PRECOMPUTE,
SURFACE_MODEL_SOURCE_FILTERS,
)
logger = logging.getLogger(__name__)
async def _get_project_notification_info(
connection: aiomysql.Connection,
project_id: UUID,
) -> dict[str, Any] | None:
async with connection.cursor(aiomysql.DictCursor) as cursor:
await cursor.execute(
"""
SELECT p.id, p.name AS project_name, u.email AS user_email, u.name AS user_name
FROM projects p
JOIN users u ON u.id = p.user_id
WHERE p.id = %s AND p.deleted_at IS NULL AND u.deleted_at IS NULL
""",
(str(project_id),),
)
row = await cursor.fetchone()
return dict(row) if row else None
async def trigger_wf1_analysis_and_email(
*,
project_id: UUID,
input_file_id: int,
user_role: str,
) -> None:
"""WF1 분석·DB 저장 후 역할에 따라 자동 확정하고 이메일을 발송한다."""
pool = get_db_pool()
project_info: dict[str, Any] | None = None
try:
source_filters = list(SURFACE_MODEL_SOURCE_FILTERS)
methods = list(SURFACE_MODEL_PRECOMPUTE)
params = {
"input_file_id": str(input_file_id),
"source_filters": source_filters,
"methods": methods,
"force": False,
}
async with pool.acquire() as connection:
async with connection.cursor() as cursor:
await start_stage(cursor, str(project_id), 1, params)
await connection.commit()
stored_path = await get_project_storage_relative_path(connection, project_id)
project_info = await _get_project_notification_info(connection, project_id)
from B04_wf1_Surface.B04_wf1_Surface_Repository import get_input_file
input_file = await get_input_file(connection, project_id, input_file_id)
project_root = Path(resolve_stored_project_path(stored_path))
las_path = project_root / Path(str(input_file["raw_file_path"]))
if not las_path.is_file():
raise FileNotFoundError("원본 LAS/LAZ 파일을 찾을 수 없습니다.")
from B04_wf1_Surface.B04_wf1_Surface_Engine import run_surface_analysis
from B04_wf1_Surface.B04_wf1_Surface_Repository import save_surface_analysis_to_db
from B04_wf1_Surface.B04_wf1_Surface_Router import write_surface_progress
write_surface_progress(project_root, 5, "analyzing", "WF1 분석을 시작합니다.")
def _on_progress(percent: int, stage: str, message: str) -> None:
write_surface_progress(project_root, percent, stage, message)
analysis_result = await asyncio.to_thread(
run_surface_analysis,
project_root,
las_path,
source_filters=source_filters,
methods=methods,
force=False,
on_progress=_on_progress,
)
auto_confirmation_error: str | None = None
auto_confirmed = False
confirmed_model_id: int | None = None
async with pool.acquire() as connection:
await connection.begin()
try:
await save_surface_analysis_to_db(
connection,
project_id=project_id,
input_file_id=input_file_id,
analysis_result=analysis_result,
source_filters=source_filters,
)
if user_role != "SYSTEM_ADMIN":
from B04_wf1_Surface.B04_wf1_Surface_Service import (
confirm_surface_selection,
find_surface_model_for_selection,
)
selection = surface_confirmation_defaults()
try:
model_id = await find_surface_model_for_selection(
connection, project_id, selection
)
except LookupError as exc:
auto_confirmation_error = str(exc)
logger.error(
"WF1 자동 확정 보류: project_id=%s reason=%s",
project_id,
auto_confirmation_error,
)
else:
await confirm_surface_selection(connection, project_id, model_id, selection)
auto_confirmed = True
confirmed_model_id = model_id
await connection.commit()
except Exception:
await connection.rollback()
raise
if auto_confirmation_error:
progress_stage = "awaiting_confirmation"
progress_message = f"WF1 자동 확정 보류 — {auto_confirmation_error}"
elif auto_confirmed:
progress_stage = "completed"
progress_message = "지표면 모델 자동 확정 완료 — 노선 설계로 이동하세요."
else:
progress_stage = "awaiting_confirmation"
progress_message = "WF1 분석이 완료되었습니다. 사용할 모델을 확정하세요."
write_surface_progress(project_root, 100, progress_stage, progress_message)
if SEND_ANALYSIS_COMPLETION_EMAIL and project_info and project_info.get("user_email"):
await send_analysis_completion_email(
project_id=project_id,
project_name=str(project_info.get("project_name") or project_id),
to_email=str(project_info["user_email"]),
analysis_result=analysis_result,
)
# WF1 자동 확정까지 끝났으면 같은 백그라운드 태스크에서 B05 기본 경로 → B06 기본
# 횡단 설계 체인을 잇는다(2026-08-04 사용자 확정). 체인은 실패를 스스로 격리한다.
if auto_confirmed and AUTO_DESIGN_CHAIN_ENABLED:
from B03_FileInput.B03_FileInput_Service_Chain import run_auto_design_chain
await run_auto_design_chain(project_id, surface_model_id=confirmed_model_id)
except Exception as exc:
logger.exception("WF1 백그라운드 분석 실패: project_id=%s", project_id)
async with pool.acquire() as connection, connection.cursor() as cursor:
await fail_stage(cursor, str(project_id), 1, str(exc))
await connection.commit()
if project_info and project_info.get("user_email"):
await send_analysis_error_email(
project_id=project_id,
project_name=str(project_info.get("project_name") or project_id),
to_email=str(project_info["user_email"]),
error_message=str(exc),
)