Files
Aislo/db_management/tools_seed_projects.py
T

233 lines
9.0 KiB
Python

"""검증 프로젝트 일괄 새로 만들기 — 같은 입력 원본으로 등록 → 입력 → 전처리 → 자동설계 체인.
복제(`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]))