diff --git a/db_management/tools_seed_projects.py b/db_management/tools_seed_projects.py new file mode 100644 index 000000000..848e130aa --- /dev/null +++ b/db_management/tools_seed_projects.py @@ -0,0 +1,232 @@ +"""검증 프로젝트 일괄 새로 만들기 — 같은 입력 원본으로 등록 → 입력 → 전처리 → 자동설계 체인. + +복제(`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]))