""" Инициализация схемы БД — идемпотентно (CREATE IF NOT EXISTS) + startup cleanup. Вызывается при ПЕРВОМ обращении к БД в каждом gunicorn-воркере. Безопасно для конкурентного вызова из нескольких воркеров — PostgreSQL корректно обрабатывает одновременный CREATE IF NOT EXISTS. ТАБЛИЦЫ: runs — история ВСЕХ запусков операций: - Ручной режим: каждая операция (create/modify/delete/...) — одна строка - Сценарный режим: каждый шаг сценария — одна строка - Колонки: id, client_id, stand, user_email, svc_id, svc_name, op_name, svc_op_id, op_uid, instance_uid, display_name, status, duration_sec, error_log, params (JSONB), stages (JSONB), app_version, scenario_run_id, step_number, instance_meta (JSONB) scenario_runs — история запусков сценариев: - Одна строка = один запуск сценария - Колонки: id, client_id, stand, scenario_name, status (RUNNING/OK/FAIL/TIMEOUT), current_step, total_steps, instance_bindings (JSONB), duration_sec, error_log, definition_id, definition_version scenario_definitions — определения сценариев (редактируются через UI): - Колонки: id, client_id, stand, name, steps (JSONB), version, is_active, created_at, updated_at, updated_by - Уникальность: (client_id, stand, LOWER(name)) — нельзя два сценария с одинаковым именем STARTUP CLEANUP: При старте приложения все зависшие RUNNING-сценарии (старше 1 часа) переводятся в TIMEOUT. Это чистит последствия падения gunicorn-воркера. SEED: При первом старте импортируются дефолтные сценарии из scenario_seed.yaml. INSERT ON CONFLICT DO NOTHING — не перезаписывает существующие. """ import os import json from db.pool import get_conn, put_conn # ── DDL для таблицы runs (история операций) ── SCHEMA_SQL = """ CREATE TABLE IF NOT EXISTS runs ( id SERIAL PRIMARY KEY, created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), client_id VARCHAR(64) NOT NULL, stand VARCHAR(16) NOT NULL, user_email VARCHAR(128), svc_id INTEGER NOT NULL, svc_name VARCHAR(128), op_name VARCHAR(32) NOT NULL, svc_op_id INTEGER, op_uid VARCHAR(64), instance_uid VARCHAR(64), display_name VARCHAR(255), status VARCHAR(16) NOT NULL, duration_sec DOUBLE PRECISION, error_log TEXT, params JSONB, stages JSONB, app_version VARCHAR(16) ); -- Индекс для фильтрации по пользователю и стенду (основной запрос истории) CREATE INDEX IF NOT EXISTS idx_runs_client_stand ON runs (client_id, stand, created_at DESC); -- Индекс для поиска по instance_uid (связь с scenario_runs) CREATE INDEX IF NOT EXISTS idx_runs_instance ON runs (instance_uid); -- Индекс по дате создания (для очистки старых записей) CREATE INDEX IF NOT EXISTS idx_runs_created ON runs (created_at); """ # ── Миграции: добавляем колонки которых может не быть в старых БД ── # ALTER TABLE ADD COLUMN IF NOT EXISTS — безопасно, не ломает существующие данные MIGRATION_SQL = """ ALTER TABLE runs ADD COLUMN IF NOT EXISTS op_uid VARCHAR(64); ALTER TABLE runs ADD COLUMN IF NOT EXISTS user_email VARCHAR(128); ALTER TABLE runs ADD COLUMN IF NOT EXISTS app_version VARCHAR(16); ALTER TABLE runs ADD COLUMN IF NOT EXISTS scenario_run_id INTEGER; ALTER TABLE runs ADD COLUMN IF NOT EXISTS step_number INTEGER; ALTER TABLE runs ADD COLUMN IF NOT EXISTS instance_meta JSONB; """ # ── DDL для scenario_runs (запуски сценариев) ── SCENARIO_RUNS_SQL = """ CREATE TABLE IF NOT EXISTS scenario_runs ( id SERIAL PRIMARY KEY, created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), client_id VARCHAR(64) NOT NULL, stand VARCHAR(16) NOT NULL, user_email VARCHAR(128), scenario_name VARCHAR(200) NOT NULL, status VARCHAR(16) NOT NULL DEFAULT 'RUNNING', current_step INTEGER NOT NULL DEFAULT 0, total_steps INTEGER NOT NULL, instance_bindings JSONB DEFAULT '{}', duration_sec DOUBLE PRECISION, error_log TEXT, app_version VARCHAR(16) ); -- Индекс для списка последних запусков (основной запрос) CREATE INDEX IF NOT EXISTS idx_scenario_runs_client_stand ON scenario_runs (client_id, stand, created_at DESC); -- Индекс для проверки lock (есть ли RUNNING) CREATE INDEX IF NOT EXISTS idx_scenario_runs_status ON scenario_runs (client_id, stand, status); -- Partial unique index — атомарный lock на уровне БД. -- Гарантирует что только ОДИН сценарий может быть RUNNING для client_id+stand. -- Вторая параллельная вставка получит unique violation → 409 без гонок. CREATE UNIQUE INDEX IF NOT EXISTS idx_one_running ON scenario_runs (client_id, stand) WHERE status = 'RUNNING'; -- Миграции для scenario_runs ALTER TABLE scenario_runs ADD COLUMN IF NOT EXISTS definition_id INTEGER; ALTER TABLE scenario_runs ADD COLUMN IF NOT EXISTS definition_version INTEGER; -- ── DDL для scenario_definitions (определения, редактируются через UI) ── CREATE TABLE IF NOT EXISTS scenario_definitions ( id SERIAL PRIMARY KEY, client_id VARCHAR(64) NOT NULL, stand VARCHAR(16) NOT NULL, name VARCHAR(200) NOT NULL, steps JSONB NOT NULL DEFAULT '[]', version INTEGER NOT NULL DEFAULT 1, is_active BOOLEAN NOT NULL DEFAULT TRUE, created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), updated_by VARCHAR(128) ); -- Уникальность имени в рамках пользователя и стенда CREATE UNIQUE INDEX IF NOT EXISTS idx_scenario_defs_unique ON scenario_definitions (client_id, stand, LOWER(name)); -- Индекс для фильтрации активных CREATE INDEX IF NOT EXISTS idx_scenario_defs_active ON scenario_definitions (client_id, stand, is_active); """ def _seed_scenarios(): """Импорт дефолтных сценариев из scenario_seed.yaml при первом старте. Формат YAML: scenarios: - name: "dummy_test" client_id: "" # пустой = seed (доступен всем) stand: "" # пустой = seed steps: [...] INSERT ON CONFLICT DO NOTHING — если пользователь уже создал сценарий с таким именем, seed не перезапишет его.""" try: import yaml, os path = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "scenario_seed.yaml") if not os.path.exists(path): return with open(path) as f: data = yaml.safe_load(f) or {} seed_list = data.get("scenarios", []) if not seed_list: return conn = get_conn() if not conn: return cur = conn.cursor() for s in seed_list: cur.execute(""" INSERT INTO scenario_definitions (client_id, stand, name, steps) VALUES (%s, %s, %s, %s) ON CONFLICT (client_id, stand, LOWER(name)) DO NOTHING """, (s.get("client_id", ""), s.get("stand", ""), s["name"], json.dumps(s.get("steps", [])))) conn.commit() cur.close() put_conn(conn) print("[DB] Seed scenarios OK", flush=True) except Exception as e: print(f"[DB] Seed scenarios error: {e}", flush=True) def init_db(): """Применить схему + миграции — идемпотентно, безопасно для конкуретного вызова. Порядок: 1. Создать таблицы (CREATE IF NOT EXISTS) 2. Применить миграции (ALTER TABLE ADD COLUMN IF NOT EXISTS) 3. Startup cleanup: зависшие RUNNING → TIMEOUT 4. Seed дефолтных сценариев """ # Без DB_USER приложение работает без БД (без истории) if not os.getenv("DB_USER"): print("[DB] DB_USER not set — skipping init", flush=True) return conn = get_conn() if not conn: print("[DB] Failed to get connection — skipping init", flush=True) return try: cur = conn.cursor() # Шаг 1-2: схема + миграции cur.execute(SCHEMA_SQL) cur.execute(MIGRATION_SQL) cur.execute(SCENARIO_RUNS_SQL) conn.commit() cur.close() print("[DB] Schema initialized", flush=True) # Шаг 3: Startup cleanup — перевести зависшие RUNNING в TIMEOUT # Если gunicorn упал во время выполнения сценария, статус остался RUNNING. # Чистим всё что старше 1 часа — операция точно не могла длиться дольше. try: cur = conn.cursor() cur.execute(""" UPDATE scenario_runs SET status = 'TIMEOUT', error_log = 'worker restart' WHERE status = 'RUNNING' AND created_at < NOW() - INTERVAL '1 hour' """) if cur.rowcount: print(f"[DB] Cleaned up {cur.rowcount} stuck scenario_runs", flush=True) conn.commit() cur.close() except Exception as e: print(f"[DB] Cleanup error: {e}", flush=True) conn.rollback() except Exception as e: print(f"[DB] Init error: {e}", flush=True) conn.rollback() finally: put_conn(conn) # Шаг 4: Seed сценариев (вне основной транзакции) _seed_scenarios()