SYSTEM_ADMIN 예외를 제거한다(2026-08-04 사용자 확정 — 관리자도 같은 사용자). 업로드 후 모두 기본값으로 지표면 자동 확정 -> B05 기본 경로 -> B06 기본 횡단 설계까지 체인을 탄다. 관리자가 B04에서 다른 모델로 재확정하는 경로는 그대로이며, 체인은 경로가 이미 있으면 스스로 건너뛴다. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
185 lines
7.8 KiB
Python
185 lines
7.8 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 저장 후 자동 확정하고 이메일을 발송한다.
|
|
|
|
자동 확정은 역할 무관이다(2026-08-04 사용자 확정 — 관리자도 같은 사용자).
|
|
`user_role`은 호출부 시그니처 유지용으로만 남겨 둔다.
|
|
"""
|
|
_ = user_role
|
|
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,
|
|
)
|
|
# 역할 무관 자동 확정(2026-08-04 사용자 확정) — 시스템 관리자도 같은
|
|
# 사용자다. 모두 기본값으로 자동 확정하고 B05·B06 자동 계산 체인까지 탄다.
|
|
# 관리자는 B04에서 다른 모델을 골라 재확정할 수 있다(체인은 경로가 이미
|
|
# 있으면 스스로 건너뛴다).
|
|
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),
|
|
)
|