1075 lines
39 KiB
Python
1075 lines
39 KiB
Python
import os
|
|
import asyncio
|
|
import datetime
|
|
import json
|
|
import locale
|
|
import sqlite3
|
|
import tempfile
|
|
import threading
|
|
import traceback
|
|
import uuid
|
|
from fastapi import APIRouter, Depends, File, HTTPException, Request, UploadFile
|
|
from sqlalchemy.orm import Session
|
|
from sqlalchemy import func, text
|
|
from pydantic import BaseModel
|
|
from typing import Optional
|
|
import io
|
|
from openpyxl import Workbook
|
|
from fastapi.responses import FileResponse, Response, StreamingResponse
|
|
from starlette.background import BackgroundTask
|
|
from urllib.parse import quote
|
|
from types import SimpleNamespace
|
|
|
|
from app.database import SessionLocal, engine, get_db
|
|
from app.database_maintenance import ensure_search_optimizations
|
|
from app.models import Order
|
|
from app.auth import get_admin_user, get_current_user
|
|
from app.event_logger import EVENT_LOG_PATH, log_event
|
|
|
|
router = APIRouter()
|
|
delete_progress_state = {}
|
|
export_progress_state = {}
|
|
export_progress_lock = threading.Lock()
|
|
DELETE_BATCH_SIZE = 50000
|
|
DELETE_COMMIT_EVERY_ROWS = 200000
|
|
|
|
# Setup correct locale for Korean day of week (e.g., 월요일)
|
|
try:
|
|
locale.setlocale(locale.LC_TIME, 'ko_KR.UTF-8')
|
|
except locale.Error:
|
|
pass
|
|
|
|
class DeleteAllOrdersRequest(BaseModel):
|
|
password: str
|
|
|
|
class DeleteOrdersByDateRangeRequest(BaseModel):
|
|
password: str
|
|
start_date: str
|
|
end_date: str
|
|
delete_token: Optional[str] = None
|
|
|
|
|
|
class UiEventRequest(BaseModel):
|
|
action: str
|
|
details: Optional[dict] = None
|
|
|
|
|
|
def format_file_size(size_bytes: int) -> str:
|
|
value = float(size_bytes or 0)
|
|
for unit in ("B", "KB", "MB", "GB", "TB"):
|
|
if value < 1024 or unit == "TB":
|
|
if unit == "B":
|
|
return f"{int(value)} {unit}"
|
|
return f"{value:.1f} {unit}"
|
|
value /= 1024
|
|
|
|
|
|
def set_delete_progress(
|
|
user_id: str,
|
|
*,
|
|
status: str,
|
|
message: str,
|
|
progress: int = 0,
|
|
cleanup_progress: int = 0,
|
|
total_rows: int = 0,
|
|
deleted_rows: int = 0,
|
|
cleanup_step: str = "",
|
|
):
|
|
delete_progress_state[user_id] = {
|
|
"status": status,
|
|
"message": message,
|
|
"progress": max(0, min(100, int(progress))),
|
|
"cleanup_progress": max(0, min(100, int(cleanup_progress))),
|
|
"total_rows": int(total_rows or 0),
|
|
"deleted_rows": int(deleted_rows or 0),
|
|
"cleanup_step": cleanup_step,
|
|
}
|
|
|
|
def require_admin_session(request: Request) -> str:
|
|
user_email = request.session.get("user_email")
|
|
if not user_email:
|
|
raise HTTPException(status_code=401, detail="Not authenticated")
|
|
|
|
db = SessionLocal()
|
|
try:
|
|
row = db.execute(
|
|
text("SELECT is_admin FROM users WHERE email = :email"),
|
|
{"email": user_email},
|
|
).fetchone()
|
|
finally:
|
|
db.close()
|
|
|
|
if not row or not row[0]:
|
|
raise HTTPException(status_code=403, detail="Admin access required")
|
|
return user_email
|
|
|
|
|
|
def excel_response_from_rows(rows, filename: str, sheet_name: str):
|
|
output = io.BytesIO()
|
|
workbook = Workbook(write_only=True)
|
|
worksheet = workbook.create_sheet(title=sheet_name[:31] or "Sheet1")
|
|
|
|
row_iter = iter(rows)
|
|
first_row = next(row_iter, None)
|
|
if first_row is None:
|
|
raise HTTPException(status_code=404, detail="다운로드할 데이터가 없습니다.")
|
|
|
|
headers = list(first_row.keys())
|
|
worksheet.append(headers)
|
|
worksheet.append([first_row.get(header) for header in headers])
|
|
for row in row_iter:
|
|
worksheet.append([row.get(header) for header in headers])
|
|
|
|
workbook.save(output)
|
|
output.seek(0)
|
|
encoded_filename = quote(filename)
|
|
return Response(
|
|
content=output.getvalue(),
|
|
media_type="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
|
|
headers={"Content-Disposition": f"attachment; filename*=UTF-8''{encoded_filename}"},
|
|
)
|
|
|
|
|
|
def write_excel_file_from_rows(rows, output_path: str, sheet_name: str):
|
|
workbook = Workbook(write_only=True)
|
|
worksheet = workbook.create_sheet(title=sheet_name[:31] or "Sheet1")
|
|
|
|
row_iter = iter(rows)
|
|
first_row = next(row_iter, None)
|
|
if first_row is None:
|
|
raise ValueError("다운로드할 데이터가 없습니다.")
|
|
|
|
headers = list(first_row.keys())
|
|
worksheet.append(headers)
|
|
worksheet.append([first_row.get(header) for header in headers])
|
|
for row in row_iter:
|
|
worksheet.append([row.get(header) for header in headers])
|
|
|
|
workbook.save(output_path)
|
|
|
|
|
|
def validate_date_range(start_date: str, end_date: str):
|
|
try:
|
|
start = datetime.datetime.strptime(start_date, "%Y-%m-%d").date()
|
|
end = datetime.datetime.strptime(end_date, "%Y-%m-%d").date()
|
|
except ValueError:
|
|
raise HTTPException(status_code=400, detail="날짜 형식이 올바르지 않습니다.")
|
|
|
|
if start > end:
|
|
raise HTTPException(status_code=400, detail="시작일이 종료일보다 늦을 수 없습니다.")
|
|
return start, end
|
|
|
|
|
|
def log_operation_error(user, action: str, exc: Exception, details: dict | None = None):
|
|
error_details = {
|
|
**(details or {}),
|
|
"error_type": type(exc).__name__,
|
|
"error": str(exc),
|
|
"traceback": traceback.format_exc(),
|
|
}
|
|
try:
|
|
log_event(user, action, "error", error_details)
|
|
except Exception:
|
|
os.makedirs("logs", exist_ok=True)
|
|
timestamp = datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
|
with open(os.path.join("logs", "operation_errors.log"), "a", encoding="utf-8") as f:
|
|
f.write(f"[{timestamp}] {action}\n")
|
|
f.write(json.dumps(error_details, ensure_ascii=False, indent=2))
|
|
f.write("\n\n")
|
|
|
|
|
|
def set_export_progress(export_token: str, **updates):
|
|
with export_progress_lock:
|
|
current = export_progress_state.get(export_token, {})
|
|
current.update(updates)
|
|
export_progress_state[export_token] = current
|
|
|
|
|
|
def get_export_progress(export_token: str):
|
|
with export_progress_lock:
|
|
job = export_progress_state.get(export_token)
|
|
return dict(job) if job else None
|
|
|
|
|
|
def build_export_rows(db: Session, start_date: Optional[str], end_date: Optional[str]):
|
|
query = db.query(
|
|
Order.order_date,
|
|
Order.recipient_name,
|
|
Order.product_code,
|
|
Order.product_name,
|
|
Order.order_quantity,
|
|
Order.address,
|
|
Order.recipient_phone,
|
|
Order.recipient_mobile,
|
|
Order.tracking_number,
|
|
Order.vendor,
|
|
Order.upload_date,
|
|
Order.order_no_mall,
|
|
Order.order_list_1,
|
|
)
|
|
|
|
if start_date and end_date:
|
|
start_datetime = start_date
|
|
end_datetime = f"{end_date} 23:59:59"
|
|
query = query.filter(Order.order_date >= start_datetime, Order.order_date <= end_datetime)
|
|
|
|
for o in query.order_by(Order.order_date.desc(), Order.id.desc()).yield_per(5000):
|
|
yield {
|
|
"주문날짜": o.order_date,
|
|
"수령인명": o.recipient_name,
|
|
"상품코드": o.product_code,
|
|
"상품명": o.product_name,
|
|
"수량": o.order_quantity,
|
|
"주소": o.address,
|
|
"수령인 전화번호": o.recipient_phone,
|
|
"휴대폰": o.recipient_mobile,
|
|
"송장번호": o.tracking_number,
|
|
"발주처": o.vendor,
|
|
"추가날짜": o.upload_date,
|
|
"주문번호(쇼핑몰)": o.order_no_mall,
|
|
"주문목록": o.order_list_1,
|
|
}
|
|
|
|
|
|
def count_export_rows(db: Session, start_date: Optional[str], end_date: Optional[str]) -> int:
|
|
query = db.query(func.count(Order.id))
|
|
if start_date and end_date:
|
|
start_datetime = start_date
|
|
end_datetime = f"{end_date} 23:59:59"
|
|
query = query.filter(Order.order_date >= start_datetime, Order.order_date <= end_datetime)
|
|
return int(query.scalar() or 0)
|
|
|
|
|
|
def run_export_job(export_token: str, start_date: Optional[str], end_date: Optional[str], user_label_data):
|
|
db = SessionLocal()
|
|
started_at = datetime.datetime.now()
|
|
export_details = {"start_date": start_date, "end_date": end_date, "export_token": export_token}
|
|
current_user = SimpleNamespace(**user_label_data)
|
|
output_path = ""
|
|
row_count = 0
|
|
|
|
try:
|
|
set_export_progress(
|
|
export_token,
|
|
status="running",
|
|
message="대상 데이터 건수를 확인 중입니다.",
|
|
processed_rows=0,
|
|
total_rows=0,
|
|
progress=0,
|
|
)
|
|
log_event(current_user, "DB 데이터 엑셀 다운로드 파일 생성 시작", "info", export_details)
|
|
|
|
total_rows = count_export_rows(db, start_date, end_date)
|
|
if total_rows == 0:
|
|
raise ValueError("해당 기간의 데이터가 없습니다.")
|
|
set_export_progress(
|
|
export_token,
|
|
message=f"총 {total_rows:,}건 엑셀 파일 생성 중입니다.",
|
|
processed_rows=0,
|
|
total_rows=total_rows,
|
|
progress=0,
|
|
)
|
|
|
|
filename = "OMS_DB_전체데이터.xlsx"
|
|
if start_date and end_date:
|
|
filename = f"OMS_DB_{start_date}_to_{end_date}.xlsx"
|
|
|
|
temp_file = tempfile.NamedTemporaryFile(prefix="order_export_", suffix=".xlsx", delete=False)
|
|
output_path = temp_file.name
|
|
temp_file.close()
|
|
|
|
workbook = Workbook(write_only=True)
|
|
worksheet = workbook.create_sheet(title="DB_Export")
|
|
headers = [
|
|
"주문날짜",
|
|
"수령인명",
|
|
"상품코드",
|
|
"상품명",
|
|
"수량",
|
|
"주소",
|
|
"수령인 전화번호",
|
|
"휴대폰",
|
|
"송장번호",
|
|
"발주처",
|
|
"추가날짜",
|
|
"주문번호(쇼핑몰)",
|
|
"주문목록",
|
|
]
|
|
worksheet.append(headers)
|
|
|
|
for row in build_export_rows(db, start_date, end_date):
|
|
row_count += 1
|
|
worksheet.append([row.get(header) for header in headers])
|
|
if row_count % 5000 == 0:
|
|
progress = round((row_count / total_rows) * 100, 1) if total_rows else 0
|
|
set_export_progress(
|
|
export_token,
|
|
message=f"{row_count:,}건 처리 중입니다.",
|
|
processed_rows=row_count,
|
|
total_rows=total_rows,
|
|
progress=min(progress, 99),
|
|
)
|
|
|
|
set_export_progress(
|
|
export_token,
|
|
message=f"{row_count:,}건 엑셀 파일 저장 중입니다.",
|
|
processed_rows=row_count,
|
|
total_rows=total_rows,
|
|
progress=99,
|
|
)
|
|
workbook.save(output_path)
|
|
|
|
elapsed_seconds = round((datetime.datetime.now() - started_at).total_seconds(), 3)
|
|
set_export_progress(
|
|
export_token,
|
|
status="completed",
|
|
message="엑셀 파일 생성이 완료되었습니다.",
|
|
processed_rows=row_count,
|
|
total_rows=total_rows,
|
|
progress=100,
|
|
file_path=output_path,
|
|
filename=filename,
|
|
elapsed_seconds=elapsed_seconds,
|
|
)
|
|
log_event(
|
|
current_user,
|
|
"DB 데이터 엑셀 다운로드 파일 생성 완료",
|
|
"success",
|
|
{
|
|
**export_details,
|
|
"download_count": row_count,
|
|
"elapsed_seconds": elapsed_seconds,
|
|
"file_path": output_path,
|
|
},
|
|
)
|
|
except Exception as e:
|
|
if output_path and os.path.exists(output_path):
|
|
try:
|
|
os.remove(output_path)
|
|
except OSError:
|
|
pass
|
|
elapsed_seconds = round((datetime.datetime.now() - started_at).total_seconds(), 3)
|
|
set_export_progress(
|
|
export_token,
|
|
status="error",
|
|
message=f"엑셀 파일 생성 오류: {type(e).__name__}: {str(e)}",
|
|
processed_rows=row_count,
|
|
total_rows=get_export_progress(export_token).get("total_rows", 0) if get_export_progress(export_token) else 0,
|
|
progress=0,
|
|
elapsed_seconds=elapsed_seconds,
|
|
)
|
|
log_operation_error(
|
|
current_user,
|
|
"DB 데이터 엑셀 다운로드 파일 생성",
|
|
e,
|
|
{
|
|
**export_details,
|
|
"processed_rows": row_count,
|
|
"elapsed_seconds": elapsed_seconds,
|
|
},
|
|
)
|
|
finally:
|
|
db.close()
|
|
|
|
|
|
@router.post("/log-ui-event")
|
|
async def log_ui_event(data: UiEventRequest, current_user = Depends(get_current_user)):
|
|
log_event(current_user, data.action, "info", data.details or {})
|
|
return {"message": "logged"}
|
|
|
|
@router.get("/test-ping")
|
|
async def test_ping():
|
|
return {"status": "ok", "message": "pong"}
|
|
|
|
@router.get("/database-size")
|
|
async def get_database_size(db: Session = Depends(get_db), current_user = Depends(get_current_user)):
|
|
try:
|
|
dialect = db.get_bind().dialect.name
|
|
size_bytes = 0
|
|
if dialect == "postgresql":
|
|
try:
|
|
size_bytes = int(db.execute(text("SELECT pg_database_size(current_database())")).scalar() or 0)
|
|
except Exception:
|
|
try:
|
|
size_bytes = int(db.execute(text("SELECT pg_total_relation_size('orders')")).scalar() or 0)
|
|
except Exception:
|
|
size_bytes = 0
|
|
elif dialect == "sqlite":
|
|
database_path = db.get_bind().url.database or "app.db"
|
|
if not os.path.isabs(database_path):
|
|
database_path = os.path.abspath(database_path)
|
|
size_bytes = os.path.getsize(database_path) if os.path.exists(database_path) else 0
|
|
else:
|
|
raise HTTPException(status_code=400, detail=f"지원하지 않는 DB 종류입니다: {dialect}")
|
|
|
|
return {
|
|
"bytes": size_bytes,
|
|
"formatted": format_file_size(size_bytes),
|
|
"database": dialect,
|
|
}
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
raise HTTPException(status_code=500, detail=f"DB 용량 조회 오류: {str(e)}")
|
|
|
|
|
|
@router.post("/optimize")
|
|
async def optimize_database(db: Session = Depends(get_db), current_user = Depends(get_current_user)):
|
|
try:
|
|
if db.get_bind().dialect.name == "sqlite":
|
|
with db.get_bind().connect() as conn:
|
|
conn.execution_options(isolation_level="AUTOCOMMIT").execute(text("VACUUM"))
|
|
else:
|
|
with db.get_bind().connect() as conn:
|
|
conn.execution_options(isolation_level="AUTOCOMMIT").execute(text("VACUUM ANALYZE orders"))
|
|
log_event(current_user, "데이터베이스 최적화", "success")
|
|
return {"message": "데이터베이스 최적화가 성공적으로 완료되었습니다."}
|
|
except Exception as e:
|
|
log_event(current_user, "데이터베이스 최적화", "error", str(e))
|
|
raise HTTPException(status_code=500, detail=f"최적화 오류: {str(e)}")
|
|
|
|
@router.get("/backup-db")
|
|
async def backup_database(current_user = Depends(get_admin_user)):
|
|
if engine.dialect.name != "sqlite":
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail="이 백업 기능은 SQLite 전용입니다. PostgreSQL은 서버 측 pg_dump를 사용하세요.",
|
|
)
|
|
backup_file = tempfile.NamedTemporaryFile(prefix="oms-backup-", suffix=".db", delete=False)
|
|
backup_path = backup_file.name
|
|
backup_file.close()
|
|
|
|
try:
|
|
source = sqlite3.connect("app.db")
|
|
target = sqlite3.connect(backup_path)
|
|
with target:
|
|
source.backup(target)
|
|
target.close()
|
|
source.close()
|
|
|
|
timestamp = datetime.datetime.now().strftime("%Y%m%d_%H%M%S")
|
|
filename = f"OMS_DB_backup_{timestamp}.db"
|
|
log_event(current_user, "DB 백업 다운로드", "success", {"filename": filename})
|
|
return FileResponse(
|
|
backup_path,
|
|
media_type="application/octet-stream",
|
|
filename=filename,
|
|
background=BackgroundTask(lambda: os.path.exists(backup_path) and os.remove(backup_path)),
|
|
)
|
|
except Exception as e:
|
|
if os.path.exists(backup_path):
|
|
os.remove(backup_path)
|
|
log_event(current_user, "DB 백업 다운로드", "error", str(e))
|
|
raise HTTPException(status_code=500, detail=f"백업 오류: {str(e)}")
|
|
|
|
@router.post("/restore-db")
|
|
async def restore_database(
|
|
request: Request,
|
|
file: UploadFile = File(...),
|
|
current_user = Depends(get_admin_user)
|
|
):
|
|
if engine.dialect.name != "sqlite":
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail="이 복구 기능은 SQLite 전용입니다. PostgreSQL은 서버 측 pg_restore를 사용하세요.",
|
|
)
|
|
filename = (file.filename or "").lower()
|
|
if not filename.endswith((".db", ".sqlite", ".sqlite3")):
|
|
raise HTTPException(status_code=400, detail="SQLite DB 백업 파일만 업로드할 수 있습니다.")
|
|
|
|
upload_file = tempfile.NamedTemporaryFile(prefix="oms-restore-upload-", suffix=".db", delete=False)
|
|
upload_path = upload_file.name
|
|
upload_file.close()
|
|
|
|
try:
|
|
contents = await file.read()
|
|
with open(upload_path, "wb") as f:
|
|
f.write(contents)
|
|
|
|
uploaded = sqlite3.connect(upload_path)
|
|
try:
|
|
try:
|
|
check_result = uploaded.execute("PRAGMA quick_check").fetchone()[0]
|
|
if check_result != "ok":
|
|
raise HTTPException(status_code=400, detail=f"DB 파일 무결성 오류: {check_result}")
|
|
|
|
required_tables = {"orders", "users"}
|
|
found_tables = {
|
|
row[0]
|
|
for row in uploaded.execute(
|
|
"SELECT name FROM sqlite_master WHERE type = 'table'"
|
|
).fetchall()
|
|
}
|
|
missing_tables = sorted(required_tables - found_tables)
|
|
if missing_tables:
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail=f"필수 테이블이 없는 DB 파일입니다: {', '.join(missing_tables)}",
|
|
)
|
|
restored_count = uploaded.execute("SELECT count(*) FROM orders").fetchone()[0]
|
|
except sqlite3.DatabaseError as e:
|
|
raise HTTPException(status_code=400, detail=f"올바른 SQLite DB 파일이 아닙니다: {str(e)}")
|
|
finally:
|
|
uploaded.close()
|
|
|
|
timestamp = datetime.datetime.now().strftime("%Y%m%d_%H%M%S")
|
|
safety_backup_path = f"app.db.pre-restore-{timestamp}.db"
|
|
|
|
current_source = sqlite3.connect("app.db")
|
|
safety_target = sqlite3.connect(safety_backup_path)
|
|
try:
|
|
with safety_target:
|
|
current_source.backup(safety_target)
|
|
finally:
|
|
safety_target.close()
|
|
current_source.close()
|
|
|
|
engine.dispose()
|
|
source = sqlite3.connect(upload_path)
|
|
target = sqlite3.connect("app.db")
|
|
try:
|
|
with target:
|
|
source.backup(target)
|
|
finally:
|
|
target.close()
|
|
source.close()
|
|
|
|
ensure_search_optimizations(engine)
|
|
log_event(
|
|
current_user,
|
|
"DB 복구",
|
|
"success",
|
|
{
|
|
"filename": file.filename,
|
|
"restored_count": restored_count,
|
|
"safety_backup": safety_backup_path,
|
|
},
|
|
)
|
|
return {
|
|
"message": "DB 복구가 완료되었습니다.",
|
|
"restored_count": restored_count,
|
|
"safety_backup": safety_backup_path,
|
|
}
|
|
except HTTPException:
|
|
log_event(current_user, "DB 복구", "error", {"filename": file.filename, "error": "invalid db file"})
|
|
raise
|
|
except Exception as e:
|
|
log_event(current_user, "DB 복구", "error", {"filename": file.filename, "error": str(e)})
|
|
raise HTTPException(status_code=500, detail=f"DB 복구 오류: {str(e)}")
|
|
finally:
|
|
if os.path.exists(upload_path):
|
|
os.remove(upload_path)
|
|
|
|
@router.post("/delete-all-orders")
|
|
async def delete_all_orders(
|
|
data: DeleteAllOrdersRequest,
|
|
db: Session = Depends(get_db),
|
|
current_user = Depends(get_admin_user)
|
|
):
|
|
if data.password != "1225":
|
|
raise HTTPException(status_code=403, detail="비밀번호가 올바르지 않습니다.")
|
|
|
|
try:
|
|
is_sqlite = db.get_bind().dialect.name == "sqlite"
|
|
if is_sqlite:
|
|
db.execute(text("DROP TRIGGER IF EXISTS orders_fts_ai"))
|
|
db.execute(text("DROP TRIGGER IF EXISTS orders_fts_ad"))
|
|
db.execute(text("DROP TRIGGER IF EXISTS orders_fts_au"))
|
|
db.execute(text("DROP TRIGGER IF EXISTS orders_phone_fts_ai"))
|
|
db.execute(text("DROP TRIGGER IF EXISTS orders_phone_fts_ad"))
|
|
db.execute(text("DROP TRIGGER IF EXISTS orders_phone_fts_au"))
|
|
result = db.execute(text("DELETE FROM orders"))
|
|
if is_sqlite:
|
|
db.execute(text("INSERT INTO orders_name_address_fts(orders_name_address_fts) VALUES('rebuild')"))
|
|
db.execute(text("INSERT INTO orders_phone_fts(orders_phone_fts) VALUES('rebuild')"))
|
|
db.commit()
|
|
ensure_search_optimizations(engine)
|
|
log_event(current_user, "전체 주문 데이터 삭제", "success", {"deleted_count": result.rowcount})
|
|
return {"message": "전체 주문 데이터가 삭제되었습니다.", "deleted_count": result.rowcount}
|
|
except Exception as e:
|
|
db.rollback()
|
|
ensure_search_optimizations(engine)
|
|
log_event(current_user, "전체 주문 데이터 삭제", "error", str(e))
|
|
raise HTTPException(status_code=500, detail=f"전체 삭제 오류: {str(e)}")
|
|
|
|
@router.get("/delete-progress")
|
|
async def delete_progress(
|
|
delete_token: str,
|
|
current_user = Depends(get_admin_user)
|
|
):
|
|
user_id = f"{current_user.id}:{delete_token}"
|
|
|
|
async def event_generator():
|
|
while True:
|
|
state = delete_progress_state.get(user_id)
|
|
if state:
|
|
yield f"data: {json.dumps(state, ensure_ascii=False)}\n\n"
|
|
if state["status"] in ["completed", "error"]:
|
|
delete_progress_state.pop(user_id, None)
|
|
break
|
|
else:
|
|
yield f"data: {json.dumps({'status': 'idle'}, ensure_ascii=False)}\n\n"
|
|
|
|
await asyncio.sleep(0.5)
|
|
|
|
return StreamingResponse(event_generator(), media_type="text/event-stream")
|
|
|
|
@router.post("/delete-orders-by-date-range")
|
|
def delete_orders_by_date_range(
|
|
data: DeleteOrdersByDateRangeRequest,
|
|
db: Session = Depends(get_db),
|
|
current_user = Depends(get_admin_user)
|
|
):
|
|
if data.password != "1225":
|
|
raise HTTPException(status_code=403, detail="비밀번호가 올바르지 않습니다.")
|
|
|
|
start_date, end_date = validate_date_range(data.start_date, data.end_date)
|
|
|
|
start_datetime = start_date.strftime("%Y-%m-%d")
|
|
end_datetime = f"{end_date.strftime('%Y-%m-%d')} 23:59:59"
|
|
|
|
progress_key = f"{current_user.id}:{data.delete_token}" if data.delete_token else None
|
|
params = {
|
|
"start_datetime": start_datetime,
|
|
"end_datetime": end_datetime,
|
|
}
|
|
|
|
try:
|
|
total_rows = db.execute(
|
|
text(
|
|
"SELECT COUNT(*) FROM orders "
|
|
"WHERE order_date >= :start_datetime "
|
|
"AND order_date <= :end_datetime"
|
|
),
|
|
params,
|
|
).scalar() or 0
|
|
|
|
if progress_key:
|
|
set_delete_progress(
|
|
progress_key,
|
|
status="processing",
|
|
message="삭제 대상을 확인했습니다.",
|
|
progress=0 if total_rows else 100,
|
|
cleanup_progress=0,
|
|
total_rows=total_rows,
|
|
deleted_rows=0,
|
|
)
|
|
|
|
if total_rows == 0:
|
|
if progress_key:
|
|
set_delete_progress(
|
|
progress_key,
|
|
status="completed",
|
|
message="삭제할 데이터가 없습니다.",
|
|
progress=100,
|
|
cleanup_progress=100,
|
|
total_rows=0,
|
|
deleted_rows=0,
|
|
)
|
|
log_event(
|
|
current_user,
|
|
"날짜 범위 주문 데이터 삭제",
|
|
"success",
|
|
{"start_date": data.start_date, "end_date": data.end_date, "deleted_count": 0},
|
|
)
|
|
return {
|
|
"message": "삭제할 주문 데이터가 없습니다.",
|
|
"deleted_count": 0,
|
|
"start_date": data.start_date,
|
|
"end_date": data.end_date,
|
|
}
|
|
|
|
if progress_key:
|
|
set_delete_progress(
|
|
progress_key,
|
|
status="processing",
|
|
message="삭제 준비 중입니다.",
|
|
progress=0,
|
|
cleanup_progress=0,
|
|
total_rows=total_rows,
|
|
deleted_rows=0,
|
|
cleanup_step="검색 인덱스 트리거 중지",
|
|
)
|
|
|
|
is_sqlite = db.get_bind().dialect.name == "sqlite"
|
|
if is_sqlite:
|
|
db.execute(text("PRAGMA temp_store = MEMORY"))
|
|
db.execute(text("DROP TRIGGER IF EXISTS orders_fts_ai"))
|
|
db.execute(text("DROP TRIGGER IF EXISTS orders_fts_ad"))
|
|
db.execute(text("DROP TRIGGER IF EXISTS orders_fts_au"))
|
|
db.execute(text("DROP TRIGGER IF EXISTS orders_phone_fts_ai"))
|
|
db.execute(text("DROP TRIGGER IF EXISTS orders_phone_fts_ad"))
|
|
db.execute(text("DROP TRIGGER IF EXISTS orders_phone_fts_au"))
|
|
db.commit()
|
|
|
|
deleted_count = 0
|
|
uncommitted_rows = 0
|
|
while deleted_count < total_rows:
|
|
result = db.execute(
|
|
text(
|
|
"DELETE FROM orders "
|
|
"WHERE id IN ("
|
|
"SELECT id FROM orders "
|
|
"WHERE order_date >= :start_datetime "
|
|
"AND order_date <= :end_datetime "
|
|
"LIMIT :batch_size"
|
|
")"
|
|
),
|
|
{
|
|
**params,
|
|
"batch_size": DELETE_BATCH_SIZE,
|
|
},
|
|
)
|
|
|
|
batch_deleted = max(0, result.rowcount or 0)
|
|
if batch_deleted == 0:
|
|
break
|
|
|
|
deleted_count += batch_deleted
|
|
uncommitted_rows += batch_deleted
|
|
if uncommitted_rows >= DELETE_COMMIT_EVERY_ROWS:
|
|
db.commit()
|
|
uncommitted_rows = 0
|
|
|
|
if progress_key:
|
|
set_delete_progress(
|
|
progress_key,
|
|
status="processing",
|
|
message="주문 데이터를 삭제하는 중입니다.",
|
|
progress=min(100, round((deleted_count / total_rows) * 100)),
|
|
cleanup_progress=0,
|
|
total_rows=total_rows,
|
|
deleted_rows=deleted_count,
|
|
)
|
|
|
|
db.commit()
|
|
|
|
if progress_key:
|
|
set_delete_progress(
|
|
progress_key,
|
|
status="processing",
|
|
message="정리 작업을 시작합니다.",
|
|
progress=100,
|
|
cleanup_progress=5,
|
|
total_rows=total_rows,
|
|
deleted_rows=deleted_count,
|
|
cleanup_step="검색 인덱스 정리 준비",
|
|
)
|
|
|
|
if progress_key:
|
|
set_delete_progress(
|
|
progress_key,
|
|
status="processing",
|
|
message="이름/주소 검색 인덱스를 정리하는 중입니다.",
|
|
progress=100,
|
|
cleanup_progress=25,
|
|
total_rows=total_rows,
|
|
deleted_rows=deleted_count,
|
|
cleanup_step="이름/주소 인덱스 정리",
|
|
)
|
|
if is_sqlite:
|
|
db.execute(text("INSERT INTO orders_name_address_fts(orders_name_address_fts) VALUES('rebuild')"))
|
|
|
|
if progress_key:
|
|
set_delete_progress(
|
|
progress_key,
|
|
status="processing",
|
|
message="전화번호 검색 인덱스를 정리하는 중입니다.",
|
|
progress=100,
|
|
cleanup_progress=60,
|
|
total_rows=total_rows,
|
|
deleted_rows=deleted_count,
|
|
cleanup_step="전화번호 인덱스 정리",
|
|
)
|
|
if is_sqlite:
|
|
db.execute(text("INSERT INTO orders_phone_fts(orders_phone_fts) VALUES('rebuild')"))
|
|
db.commit()
|
|
|
|
if progress_key:
|
|
set_delete_progress(
|
|
progress_key,
|
|
status="processing",
|
|
message="검색 기능을 다시 연결하는 중입니다.",
|
|
progress=100,
|
|
cleanup_progress=90,
|
|
total_rows=total_rows,
|
|
deleted_rows=deleted_count,
|
|
cleanup_step="검색 트리거 복구",
|
|
)
|
|
ensure_search_optimizations(engine)
|
|
if progress_key:
|
|
set_delete_progress(
|
|
progress_key,
|
|
status="completed",
|
|
message="삭제가 완료되었습니다.",
|
|
progress=100,
|
|
cleanup_progress=100,
|
|
total_rows=total_rows,
|
|
deleted_rows=deleted_count,
|
|
cleanup_step="정리 완료",
|
|
)
|
|
log_event(
|
|
current_user,
|
|
"날짜 범위 주문 데이터 삭제",
|
|
"success",
|
|
{"start_date": data.start_date, "end_date": data.end_date, "deleted_count": deleted_count},
|
|
)
|
|
return {
|
|
"message": "지정한 날짜 범위의 주문 데이터가 삭제되었습니다.",
|
|
"deleted_count": deleted_count,
|
|
"start_date": data.start_date,
|
|
"end_date": data.end_date,
|
|
}
|
|
except Exception as e:
|
|
db.rollback()
|
|
ensure_search_optimizations(engine)
|
|
if progress_key:
|
|
set_delete_progress(
|
|
progress_key,
|
|
status="error",
|
|
message=f"삭제 오류: {str(e)}",
|
|
progress=0,
|
|
cleanup_progress=0,
|
|
)
|
|
log_event(
|
|
current_user,
|
|
"날짜 범위 주문 데이터 삭제",
|
|
"error",
|
|
{"start_date": data.start_date, "end_date": data.end_date, "error": str(e)},
|
|
)
|
|
raise HTTPException(status_code=500, detail=f"날짜 범위 삭제 오류: {str(e)}")
|
|
|
|
@router.post("/export-db/start")
|
|
def start_export_db(
|
|
start_date: Optional[str] = None,
|
|
end_date: Optional[str] = None,
|
|
current_user = Depends(get_current_user)
|
|
):
|
|
if (start_date and not end_date) or (end_date and not start_date):
|
|
raise HTTPException(status_code=400, detail="시작일과 종료일을 모두 선택해주세요.")
|
|
if start_date and end_date:
|
|
validate_date_range(start_date, end_date)
|
|
|
|
export_token = str(uuid.uuid4())
|
|
user_label_data = {
|
|
"email": getattr(current_user, "email", None),
|
|
"name": getattr(current_user, "name", None),
|
|
"id": getattr(current_user, "id", ""),
|
|
}
|
|
set_export_progress(
|
|
export_token,
|
|
status="queued",
|
|
message="엑셀 파일 생성 대기 중입니다.",
|
|
processed_rows=0,
|
|
start_date=start_date,
|
|
end_date=end_date,
|
|
)
|
|
threading.Thread(
|
|
target=run_export_job,
|
|
args=(export_token, start_date, end_date, user_label_data),
|
|
daemon=True,
|
|
).start()
|
|
log_event(
|
|
current_user,
|
|
"DB 데이터 엑셀 다운로드 작업 등록",
|
|
"info",
|
|
{"start_date": start_date, "end_date": end_date, "export_token": export_token},
|
|
)
|
|
return {"export_token": export_token, "status": "queued"}
|
|
|
|
|
|
@router.get("/export-db/status")
|
|
def get_export_db_status(export_token: str, current_user = Depends(get_current_user)):
|
|
job = get_export_progress(export_token)
|
|
if not job:
|
|
raise HTTPException(status_code=404, detail="다운로드 작업을 찾을 수 없습니다.")
|
|
public_job = {k: v for k, v in job.items() if k != "file_path"}
|
|
return public_job
|
|
|
|
|
|
@router.get("/export-db/download")
|
|
def download_export_db_file(export_token: str, current_user = Depends(get_current_user)):
|
|
job = get_export_progress(export_token)
|
|
if not job:
|
|
raise HTTPException(status_code=404, detail="다운로드 작업을 찾을 수 없습니다.")
|
|
if job.get("status") != "completed":
|
|
raise HTTPException(status_code=409, detail=job.get("message") or "엑셀 파일 생성이 아직 완료되지 않았습니다.")
|
|
|
|
file_path = job.get("file_path")
|
|
if not file_path or not os.path.exists(file_path):
|
|
raise HTTPException(status_code=404, detail="생성된 엑셀 파일을 찾을 수 없습니다.")
|
|
|
|
log_event(
|
|
current_user,
|
|
"DB 데이터 엑셀 파일 다운로드",
|
|
"success",
|
|
{
|
|
"export_token": export_token,
|
|
"processed_rows": job.get("processed_rows", 0),
|
|
"filename": job.get("filename"),
|
|
},
|
|
)
|
|
return FileResponse(
|
|
file_path,
|
|
media_type="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
|
|
filename=job.get("filename") or "DB_Export.xlsx",
|
|
)
|
|
|
|
|
|
@router.get("/export-db")
|
|
def export_entire_db_legacy(current_user = Depends(get_current_user)):
|
|
raise HTTPException(
|
|
status_code=409,
|
|
detail="엑셀 파일 생성 시간이 길어 백그라운드 다운로드 방식으로 변경되었습니다. 페이지를 새로고침한 뒤 다시 시도해주세요.",
|
|
)
|
|
|
|
|
|
@router.get("/download-log")
|
|
def download_event_log(current_user = Depends(get_admin_user)):
|
|
try:
|
|
log_event(current_user, "로그 파일 다운로드 요청", "info")
|
|
timestamp = datetime.datetime.now().strftime("%Y%m%d_%H%M%S")
|
|
filename = quote(f"OMS_event_log_{timestamp}.csv")
|
|
with open(EVENT_LOG_PATH, "rb") as f:
|
|
content = f.read()
|
|
log_event(
|
|
current_user,
|
|
"로그 파일 다운로드",
|
|
"success",
|
|
{"event_log_path": os.path.abspath(EVENT_LOG_PATH), "bytes": len(content)},
|
|
)
|
|
return Response(
|
|
content=content,
|
|
media_type="text/csv; charset=utf-8",
|
|
headers={"Content-Disposition": f"attachment; filename*=UTF-8''{filename}"},
|
|
)
|
|
except Exception as e:
|
|
log_operation_error(
|
|
current_user,
|
|
"로그 파일 다운로드",
|
|
e,
|
|
{"event_log_path": os.path.abspath(EVENT_LOG_PATH)},
|
|
)
|
|
raise HTTPException(status_code=500, detail=f"로그 다운로드 오류: {type(e).__name__}: {str(e)}")
|
|
|
|
|
|
@router.get("/repeat-purchase-export")
|
|
async def export_repeat_purchase_orders(
|
|
start_date: str,
|
|
end_date: str,
|
|
min_count: int = 3,
|
|
db: Session = Depends(get_db),
|
|
current_user = Depends(get_admin_user),
|
|
):
|
|
validate_date_range(start_date, end_date)
|
|
if min_count < 1:
|
|
raise HTTPException(status_code=400, detail="구매 회수는 1 이상이어야 합니다.")
|
|
|
|
params = {
|
|
"start_datetime": start_date,
|
|
"end_datetime": f"{end_date} 23:59:59",
|
|
"min_count": min_count,
|
|
}
|
|
|
|
try:
|
|
rows = db.execute(
|
|
text(
|
|
"""
|
|
WITH qualified AS (
|
|
SELECT
|
|
trim(coalesce(recipient_name, '')) AS name_key,
|
|
coalesce(
|
|
nullif(recipient_mobile_normalized, ''),
|
|
nullif(recipient_phone_normalized, '')
|
|
) AS phone_key,
|
|
count(DISTINCT substr(order_date, 1, 10)) AS purchase_count
|
|
FROM orders
|
|
WHERE order_date >= :start_datetime
|
|
AND order_date <= :end_datetime
|
|
AND trim(coalesce(recipient_name, '')) != ''
|
|
AND coalesce(
|
|
nullif(recipient_mobile_normalized, ''),
|
|
nullif(recipient_phone_normalized, '')
|
|
) IS NOT NULL
|
|
GROUP BY name_key, phone_key
|
|
HAVING count(DISTINCT substr(order_date, 1, 10)) >= :min_count
|
|
)
|
|
SELECT
|
|
q.purchase_count AS purchase_count,
|
|
o.order_date AS order_date,
|
|
o.recipient_name AS recipient_name,
|
|
coalesce(nullif(o.recipient_mobile, ''), o.recipient_phone) AS display_phone,
|
|
o.product_code AS product_code,
|
|
o.product_name AS product_name,
|
|
o.order_quantity AS order_quantity,
|
|
o.address AS address,
|
|
o.tracking_number AS tracking_number,
|
|
o.vendor AS vendor,
|
|
o.order_no AS order_no,
|
|
o.order_no_mall AS order_no_mall,
|
|
o.order_list_1 AS order_list_1
|
|
FROM orders o
|
|
JOIN qualified q
|
|
ON trim(coalesce(o.recipient_name, '')) = q.name_key
|
|
AND coalesce(
|
|
nullif(o.recipient_mobile_normalized, ''),
|
|
nullif(o.recipient_phone_normalized, '')
|
|
) = q.phone_key
|
|
WHERE o.order_date >= :start_datetime
|
|
AND o.order_date <= :end_datetime
|
|
ORDER BY q.purchase_count DESC, q.name_key ASC, q.phone_key ASC, o.order_date DESC, o.id DESC
|
|
"""
|
|
),
|
|
params,
|
|
).mappings().all()
|
|
|
|
if not rows:
|
|
raise HTTPException(status_code=404, detail="해당 조건의 반복 구매 주문이 없습니다.")
|
|
|
|
excel_rows = []
|
|
for row in rows:
|
|
excel_rows.append({
|
|
"구매회수": row["purchase_count"],
|
|
"주문날짜": row["order_date"],
|
|
"수령인명": row["recipient_name"],
|
|
"전화번호": row["display_phone"],
|
|
"상품코드": row["product_code"],
|
|
"상품명": row["product_name"],
|
|
"수량": row["order_quantity"],
|
|
"주소": row["address"],
|
|
"송장번호": row["tracking_number"],
|
|
"발주처": row["vendor"],
|
|
"주문번호": row["order_no"],
|
|
"주문번호(쇼핑몰)": row["order_no_mall"],
|
|
"주문목록": row["order_list_1"],
|
|
})
|
|
|
|
filename = f"반복구매_{start_date}_to_{end_date}_{min_count}회이상.xlsx"
|
|
log_event(
|
|
current_user,
|
|
"반복 구매 주문 리스트 다운로드",
|
|
"success",
|
|
{
|
|
"start_date": start_date,
|
|
"end_date": end_date,
|
|
"min_count": min_count,
|
|
"download_count": len(excel_rows),
|
|
},
|
|
)
|
|
return excel_response_from_rows(excel_rows, filename, "Repeat_Purchases")
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
log_event(
|
|
current_user,
|
|
"반복 구매 주문 리스트 다운로드",
|
|
"error",
|
|
{
|
|
"start_date": start_date,
|
|
"end_date": end_date,
|
|
"min_count": min_count,
|
|
"error": str(e),
|
|
},
|
|
)
|
|
raise HTTPException(status_code=500, detail=f"반복 구매 리스트 다운로드 오류: {str(e)}")
|