- bd/__init__.py: добавлены get_engine() и get_session() — один SQLAlchemy engine на весь процесс - make_engine(): pool_size=5, max_overflow=5, pool_timeout=30, pool_recycle=1800 - Все 17 route-файлов: убраны локальные get_session()/make_engine(), импорт из bd - transfer_crud.py: _get_session() → get_session() из bd - init_data_base.py: get_engine() вместо make_engine() - Результат: 78 idle соединений → 1
174 lines
6.5 KiB
Python
174 lines
6.5 KiB
Python
from fastapi import APIRouter, Depends, HTTPException
|
||
import pkgutil
|
||
import importlib
|
||
from pathlib import Path
|
||
from typing import List
|
||
|
||
from bd import Settings, get_engine
|
||
from sqlalchemy import text
|
||
from sqlalchemy.exc import SQLAlchemyError
|
||
from pydantic import BaseModel
|
||
|
||
from route.auth_utils import require_admin_key
|
||
|
||
router = APIRouter(tags=["db"])
|
||
|
||
|
||
def _iter_table_modules() -> List[str]:
|
||
try:
|
||
import bd.tables as tables_pkg
|
||
pkg_paths = getattr(tables_pkg, "__path__", None)
|
||
if not pkg_paths:
|
||
pkg_paths = [str(Path(__file__).resolve().parent.parent / "bd" / "tables")]
|
||
except Exception:
|
||
pkg_paths = [str(Path(__file__).resolve().parent.parent / "bd" / "tables")]
|
||
|
||
names = []
|
||
for finder, name, ispkg in pkgutil.iter_modules(pkg_paths):
|
||
names.append(name)
|
||
return names
|
||
|
||
|
||
def collect_metadatas():
|
||
metadatas = []
|
||
for mod_name in _iter_table_modules():
|
||
try:
|
||
module = importlib.import_module(f"bd.tables.{mod_name}")
|
||
except Exception:
|
||
continue
|
||
Base = getattr(module, "Base", None)
|
||
if Base is not None and hasattr(Base, "metadata"):
|
||
metadatas.append(Base.metadata)
|
||
return metadatas
|
||
|
||
|
||
@router.post("/db/create-tables", dependencies=[Depends(require_admin_key)])
|
||
async def create_tables():
|
||
"""Создаёт все таблицы, описанные в модулях `db.tables`.
|
||
|
||
Endpoint вызывается по нажатию кнопки в UI (POST).
|
||
"""
|
||
engine = get_engine()
|
||
|
||
metadatas = collect_metadatas()
|
||
if not metadatas:
|
||
raise HTTPException(status_code=400, detail="No table metadata found in db.tables")
|
||
|
||
# deduplicate metadata objects (multiple modules may expose the same Base.metadata)
|
||
unique = []
|
||
seen = set()
|
||
for md in metadatas:
|
||
if id(md) not in seen:
|
||
seen.add(id(md))
|
||
unique.append(md)
|
||
|
||
try:
|
||
for md in unique:
|
||
md.create_all(bind=engine)
|
||
return {"status": "ok", "detail": f"Created {len(unique)} metadata groups"}
|
||
except SQLAlchemyError as e:
|
||
raise HTTPException(status_code=500, detail=str(e))
|
||
|
||
|
||
@router.post("/db/migrate-cascade-user", dependencies=[Depends(require_admin_key)])
|
||
async def migrate_cascade_user():
|
||
"""Добавляет FK responses.user_id → users.id ON DELETE CASCADE, если ещё не существует."""
|
||
engine = get_engine()
|
||
try:
|
||
with engine.connect() as conn:
|
||
conn.execute(text("""
|
||
DO $$ BEGIN
|
||
IF NOT EXISTS (
|
||
SELECT 1 FROM pg_constraint
|
||
WHERE conname = 'responses_user_id_fkey'
|
||
) THEN
|
||
ALTER TABLE responses
|
||
ADD CONSTRAINT responses_user_id_fkey
|
||
FOREIGN KEY (user_id) REFERENCES users(id) ON DELETE CASCADE;
|
||
END IF;
|
||
END $$;
|
||
"""))
|
||
conn.commit()
|
||
return {"status": "ok", "detail": "FK responses.user_id → users(id) ON DELETE CASCADE ensured"}
|
||
except SQLAlchemyError as e:
|
||
raise HTTPException(status_code=500, detail=str(e))
|
||
|
||
|
||
@router.post("/db/migrate-users", dependencies=[Depends(require_admin_key)])
|
||
async def migrate_users():
|
||
"""Добавляет колонки username и hashed_password в таблицу users, если они ещё не существуют."""
|
||
engine = get_engine()
|
||
try:
|
||
with engine.connect() as conn:
|
||
conn.execute(text("""
|
||
ALTER TABLE users
|
||
ADD COLUMN IF NOT EXISTS username VARCHAR(150),
|
||
ADD COLUMN IF NOT EXISTS hashed_password VARCHAR(256)
|
||
"""))
|
||
conn.commit()
|
||
return {"status": "ok", "detail": "Columns username and hashed_password ensured in users table"}
|
||
except SQLAlchemyError as e:
|
||
raise HTTPException(status_code=500, detail=str(e))
|
||
|
||
|
||
class ClearDBIn(BaseModel):
|
||
confirm: bool
|
||
|
||
|
||
@router.post("/db/migrate", dependencies=[Depends(require_admin_key)])
|
||
async def migrate_tables():
|
||
"""Приводит схему БД в соответствие с моделями: убирает устаревшие колонки, добавляет новые."""
|
||
engine = get_engine()
|
||
migrations = [
|
||
# Убираем старый FK и колонку group.user_id (если остался от прежней схемы)
|
||
'ALTER TABLE "group" DROP COLUMN IF EXISTS user_id',
|
||
# Делаем organization.group_id необязательным
|
||
"ALTER TABLE organization ALTER COLUMN group_id DROP NOT NULL",
|
||
# Добавляем новые колонки в users
|
||
'ALTER TABLE users ADD COLUMN IF NOT EXISTS group_id UUID REFERENCES "group"(id) ON DELETE SET NULL',
|
||
"ALTER TABLE users ADD COLUMN IF NOT EXISTS organization_id UUID REFERENCES organization(id) ON DELETE SET NULL",
|
||
# Добавляем путь к SVG-файлу диаграммы в radar_results
|
||
"ALTER TABLE radar_results ADD COLUMN IF NOT EXISTS image_path VARCHAR(512)",
|
||
]
|
||
applied = []
|
||
try:
|
||
with engine.begin() as conn:
|
||
for sql in migrations:
|
||
conn.execute(text(sql))
|
||
applied.append(sql)
|
||
return {"status": "ok", "applied": applied}
|
||
except SQLAlchemyError as e:
|
||
raise HTTPException(status_code=500, detail=str(e))
|
||
|
||
|
||
@router.post("/db/clear", dependencies=[Depends(require_admin_key)])
|
||
async def clear_tables(payload: ClearDBIn):
|
||
"""Полная очистка всех таблиц, описанных в `bd.tables`.
|
||
|
||
Требуется явное подтверждение: POST с телом {"confirm": true}.
|
||
"""
|
||
if not payload.confirm:
|
||
raise HTTPException(status_code=400, detail="Confirmation required")
|
||
|
||
engine = get_engine()
|
||
|
||
metadatas = collect_metadatas()
|
||
if not metadatas:
|
||
raise HTTPException(status_code=400, detail="No table metadata found in db.tables")
|
||
|
||
# deduplicate metadata objects
|
||
unique = []
|
||
seen = set()
|
||
for md in metadatas:
|
||
if id(md) not in seen:
|
||
seen.add(id(md))
|
||
unique.append(md)
|
||
|
||
try:
|
||
with engine.begin() as conn:
|
||
conn.execute(text("DROP SCHEMA public CASCADE"))
|
||
conn.execute(text("CREATE SCHEMA public"))
|
||
return {"status": "ok", "detail": "All tables dropped via DROP SCHEMA CASCADE"}
|
||
except SQLAlchemyError as e:
|
||
raise HTTPException(status_code=500, detail=str(e))
|