""" Инициализация схемы БД — идемпотентно (CREATE IF NOT EXISTS). Вызывается при старте приложения (в app.py, до register_blueprint). Безопасно для нескольких gunicorn-воркеров — PostgreSQL корректно обрабатывает конкурентный DDL. Таблица runs — история запусков операций: - id, created_at - client_id, stand (изоляция по пользователю и стенду) - user_email (кто запустил) - svc_id, svc_name, op_name, svc_op_id - instance_uid, display_name, op_uid - status (OK/FAIL/TIMEOUT), duration_sec, error_log - params (JSONB), stages (JSONB) - app_version (версия приложения) """ import os import json from db.pool import get_conn, put_conn 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); CREATE INDEX IF NOT EXISTS idx_runs_instance ON runs (instance_uid); CREATE INDEX IF NOT EXISTS idx_runs_created ON runs (created_at); """ # Миграция для существующих таблиц (добавляем колонки если их нет) 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; """ # Таблица сценариев — один запуск = одна строка 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); CREATE INDEX IF NOT EXISTS idx_scenario_runs_status ON scenario_runs (client_id, stand, status); ALTER TABLE scenario_runs ADD COLUMN IF NOT EXISTS definition_id INTEGER; ALTER TABLE scenario_runs ADD COLUMN IF NOT EXISTS definition_version INTEGER; -- Таблица определений сценариев — редактируются без редеплоя 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 при первом старте.""" 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(): """Применить схему + миграции — идемпотентно, безопасно для конкурентного вызова.""" 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() cur.execute(SCHEMA_SQL) cur.execute(MIGRATION_SQL) cur.execute(SCENARIO_RUNS_SQL) conn.commit() cur.close() print("[DB] Schema initialized", flush=True) # Startup cleanup: зависшие RUNNING сценарии старше 1 часа → TIMEOUT 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) # Seed сценариев — только при первом старте (INSERT ON CONFLICT DO NOTHING) _seed_scenarios()