Files
dbx-orderlist/app/routers/settings.py
T

1064 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("/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
if dialect == "postgresql":
size_bytes = int(db.execute(text("SELECT pg_database_size(current_database())")).scalar() or 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)}")