chore(db): 54-44 검증 프로젝트 일괄 새로 만들기 도구 — 원본 입력 복사 → 업로드 같은 등록 → WF1 · 자동설계 체인 서버 함수 직접 호출 · 메일 끔 · B05 읽기 확인
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UoSr2sqvYtBf1dtQdus88s
This commit is contained in:
@@ -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]))
|
||||
Reference in New Issue
Block a user