""" Миграционный раннер для py_service. Автоматически применяет SQL-файлы из папки migrations/ при старте приложения. Отслеживает применённые миграции в таблице _applied_migrations. Порядок: 1. Создаёт таблицу _applied_migrations (если нет) 2. Сканирует migrations/*.sql 3. Применяет неприменённые по порядку имени файла 4. Записывает имя файла и время в _applied_migrations Формат SQL-файлов: - Имя: 001_description.sql, 002_description.sql и т.д. - Содержимое: один или несколько SQL-запросов, разделённых точкой с запятой - Каждый запрос разделяется через ; (точка с запятой на отдельной строке или в конце) """ import os import glob import logging from datetime import datetime from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy import text logger = logging.getLogger("migrations") MIGRATIONS_DIR = os.path.join(os.path.dirname(os.path.dirname(__file__)), "migrations") async def ensure_tracking_table(db: AsyncSession) -> None: """Создаёт таблицу _applied_migrations если она ещё не существует.""" await db.execute(text(""" CREATE TABLE IF NOT EXISTS _applied_migrations ( id SERIAL PRIMARY KEY, filename VARCHAR(255) NOT NULL UNIQUE, applied_at TIMESTAMP NOT NULL DEFAULT NOW() ) """)) await db.commit() async def get_applied_migrations(db: AsyncSession) -> set: """Возвращает множество имён уже применённых миграций.""" result = await db.execute(text("SELECT filename FROM _applied_migrations")) return {row[0] for row in result.fetchall()} def _split_sql(content: str) -> list: """Разделяет SQL-содержимое на отдельные запросы. Использует точку с запятой как разделитель. Пустые запросы и комментарии пропускаются. """ queries = [] for part in content.split(";"): # Убираем SQL-комментарии (строки начинающиеся с --) lines = [] for line in part.split("\n"): stripped = line.strip() if stripped.startswith("--"): continue lines.append(line) cleaned = "\n".join(lines).strip() if cleaned: queries.append(cleaned) return queries async def apply_migrations(db: AsyncSession) -> list: """Применяет все неприменённые миграции из папки migrations/. Возвращает список имён применённых миграций. """ await ensure_tracking_table(db) # Получаем уже применённые applied = await get_applied_migrations(db) # Сканируем SQL-файлы pattern = os.path.join(MIGRATIONS_DIR, "*.sql") files = sorted(glob.glob(pattern)) applied_now = [] for filepath in files: filename = os.path.basename(filepath) # Пропускаем уже применённые if filename in applied: continue # Читаем SQL with open(filepath, "r", encoding="utf-8") as f: content = f.read() # Разделяем на запросы queries = _split_sql(content) if not queries: continue # Применяем каждый запрос try: for query in queries: try: await db.execute(text(query)) except Exception as query_err: # Idempotent DDL: CREATE TABLE IF NOT EXISTS, CREATE INDEX IF NOT EXISTS, # ALTER TABLE ADD COLUMN IF NOT EXISTS — пропускаем если объект уже существует err_msg = str(query_err).lower() is_idempotent = ( "already exists" in err_msg or "duplicate key value violates unique constraint" in err_msg ) if is_idempotent: logger.info( "Пропускаем (уже существует): %.100s", query[:100] ) continue raise # Записываем в tracking table await db.execute( text("INSERT INTO _applied_migrations (filename, applied_at) VALUES (:fn, NOW())"), {"fn": filename} ) await db.commit() applied_now.append(filename) except Exception as e: await db.rollback() raise RuntimeError(f"Ошибка при применении миграции {filename}: {e}") return applied_now async def run_migrations_on_startup() -> list: """Запуск миграций при старте приложения. Используется из lifespan в main.py. Возвращает список применённых миграций. """ from app.database import async_session async with async_session() as db: return await apply_migrations(db)