"""검증 프로젝트 일괄 새로 만들기 — 같은 입력 원본으로 등록 → 입력 → 전처리 → 자동설계 체인. 복제(`tools_clone_project.py`)와 다름 — 옛 결과를 뜨지 않고 **지금 로직으로 처음부터** 계산. 화면 없이 서버 함수를 직접 부름(로그인 없음 · 소유자는 등록 함수의 user_id). - 원본은 읽기만 — DB 행(등록 칸)과 `B03_FileInput/input` 파일을 복사해 씀. - 이름이 이미 있으면(삭제 안 된 행) 건너뜀 — 다시 돌려도 이어서 만듦. - 알림 메일 끔. 쓰는 법 ./venv/Scripts/python.exe db_management/tools_seed_projects.py \ <원본 UUID> <결과 md> 이름:user_id ... """ import asyncio import json import logging import os import shutil import sys import time from pathlib import Path from uuid import UUID os.environ["SEND_ANALYSIS_COMPLETION_EMAIL"] = "False" sys.path.append(str(Path(__file__).resolve().parent.parent)) import aiomysql import B03_FileInput.B03_FileInput_Service_WF1 as wf1 from B02_ProjRegister.B02_ProjRegister_Repository import create_project from B03_FileInput.B03_FileInput_Engine_Analyze import analyze_input_metadata from B03_FileInput.B03_FileInput_Repository import create_input_file from B03_FileInput.B03_FileInput_Router_Helpers import ( _complete_file_input_if_ready, _write_stage_metadata, ) from B03_FileInput.B03_FileInput_Schema import UploadedFileResult from B04_PreProcess.B04_PreProcess_Router_Basins import get_pipe_points from B05_Profile.B05_Profile_Structures_Router import read_structures from common_util.common_util_initial_snapshot import is_design_failed, read_design_failure from common_util.common_util_storage import resolve_stored_project_path from config.config_db import close_db_pool, get_db_pool, init_db_pool wf1.SEND_ANALYSIS_COMPLETION_EMAIL = False async def _no_mail(**_: object) -> None: return None wf1.send_analysis_error_email = _no_mail logger = logging.getLogger("seed") REGISTER_FIELDS = ( "region", "road_type", "project_year", "estimated_length_m", "memo", ) TITLE_FIELDS = ( "client_org", "project_number", "work_amount", "design_date", "pm_user_id", "field_lead_user_id", "designer_user_id", "logo_asset_id", "route_start_m", "route_end_m", ) async def _one(sql: str, args: tuple) -> dict | None: async with get_db_pool().acquire() as c, c.cursor(aiomysql.DictCursor) as cur: await cur.execute(sql, args) return await cur.fetchone() async def _register_inputs(project_id: str, source_root: Path, source_files: list[dict]) -> int: """원본 입력 파일을 같은 상대 경로로 복사하고 업로드 라우터와 같은 DB 등록을 한다.""" pool = get_db_pool() project_uuid = UUID(project_id) row = await _one("SELECT storage_path FROM projects WHERE id=%s", (project_id,)) project_root = Path(resolve_stored_project_path(row["storage_path"])) results: list[UploadedFileResult] = [] async with pool.acquire() as connection: await connection.begin() try: for item in source_files: relative_path = item["relative_path"] destination = project_root / relative_path destination.parent.mkdir(parents=True, exist_ok=True) await asyncio.to_thread(shutil.copyfile, source_root / relative_path, destination) metadata = await asyncio.to_thread(analyze_input_metadata, destination) size = destination.stat().st_size crs_epsg = metadata.get("epsg") input_file_id = await create_input_file( connection, project_id=project_uuid, file_type=item["file_type"], original_filename=item["original_filename"], relative_path=relative_path, file_size_bytes=size, upload_by=None, crs_epsg=int(crs_epsg) if crs_epsg is not None else None, metadata=metadata, ) results.append( UploadedFileResult( input_file_id=input_file_id, original_filename=item["original_filename"], file_type=item["file_type"], relative_path=relative_path, size_bytes=size, metadata=metadata, ) ) point_cloud_input_id = await _complete_file_input_if_ready(connection, project_uuid) await connection.commit() except Exception: await connection.rollback() raise _write_stage_metadata(project_root / "B03_FileInput", project_uuid, results) return point_cloud_input_id async def _check_b05(project_id: str) -> str: """B05 를 여는 두 읽기(구조물 · 관 지점)와 관 하나의 options.""" structures = await read_structures(UUID(project_id)) pipes = await get_pipe_points(UUID(project_id)) if not isinstance(pipes, dict) or not hasattr(structures, "structures"): return "B05 읽기 실패" points = pipes.get("pipe_points") or [] first = (points[0].get("options") if points else None) or {} return ( f"구조물 {len(structures.structures)} · 관 {len(points)} · " f"첫 관 options {json.dumps(first, ensure_ascii=False)}" ) async def _seed(name: str, user_id: int, source: dict, source_files: list[dict]) -> dict: existing = await _one( "SELECT id FROM projects WHERE name=%s AND user_id=%s AND deleted_at IS NULL", (name, user_id), ) if existing: project_id = existing["id"] return { "name": name, "user_id": user_id, "id": project_id, "note": "이미 있음", "b05": await _check_b05(project_id), } started = time.monotonic() created = await create_project( user_id=user_id, company_id=int(source["company_id"]), name=name, title_block={key: source[key] for key in TITLE_FIELDS}, **{key: source[key] for key in REGISTER_FIELDS}, ) project_id = str(created["project_id"]) logger.info("만듦 %s %s", name, project_id) input_id = await _register_inputs( project_id, Path(resolve_stored_project_path(source["storage_path"])), source_files, ) await wf1.trigger_wf1_analysis_and_email( project_id=UUID(project_id), input_file_id=input_id, user_role="", started_by=user_id ) row = await _one("SELECT status, storage_path FROM projects WHERE id=%s", (project_id,)) root = Path(resolve_stored_project_path(row["storage_path"])) failure = read_design_failure(root) if is_design_failed(root) else "" return { "name": name, "user_id": user_id, "id": project_id, "status": row["status"], "failure": failure, "snapshot": (root / "initial_snapshot").is_dir(), "b05": await _check_b05(project_id), "minutes": round((time.monotonic() - started) / 60, 1), } async def main(source_id: str, report: Path, targets: list[tuple[str, int]]) -> None: await init_db_pool() try: source = await _one("SELECT * FROM projects WHERE id=%s", (source_id,)) source_root = Path(resolve_stored_project_path(source["storage_path"])) meta = json.loads((source_root / "B03_FileInput/metadata.json").read_text("utf-8")) source_files = meta["files"] rows = [] for name, user_id in targets: try: result = await _seed(name, user_id, source, source_files) except Exception as exc: logger.exception("실패 %s", name) result = {"name": name, "user_id": user_id, "id": "", "failure": repr(exc)} logger.info("결과 %s", json.dumps(result, ensure_ascii=False)) rows.append(result) _write_report(report, source_id, rows) finally: await close_db_pool() def _write_report(report: Path, source_id: str, rows: list[dict]) -> None: lines = [ f"# 새 검증 프로젝트 (원본 {source_id} 입력 · 새 로직 처음부터)", "", "| 이름 | 소유 user | 프로젝트 id | 상태 | 초기값 | B05 확인 | 실패 | 분 |", "| --- | --- | --- | --- | --- | --- | --- | --- |", ] for r in rows: lines.append( f"| {r['name']} | {r['user_id']} | {r['id']} | {r.get('status', r.get('note', ''))} " f"| {r.get('snapshot', '')} | {r.get('b05', '')} | {r.get('failure', '')} " f"| {r.get('minutes', '')} |" ) report.parent.mkdir(parents=True, exist_ok=True) report.write_text("\n".join(lines) + "\n", encoding="utf-8") if __name__ == "__main__": logging.basicConfig(level=logging.INFO, format="%(asctime)s %(name)s %(message)s") pairs = [arg.rsplit(":", 1) for arg in sys.argv[3:]] asyncio.run(main(sys.argv[1], Path(sys.argv[2]), [(n, int(u)) for n, u in pairs]))