245 lines
11 KiB
Python
245 lines
11 KiB
Python
"""
|
||
Инициализация схемы БД — идемпотентно (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()
|