Рефакторинг API-слоя, операций и БД
This commit is contained in:
+71
-20
@@ -1,25 +1,45 @@
|
||||
"""
|
||||
Инициализация схемы БД — идемпотентно (CREATE IF NOT EXISTS).
|
||||
Инициализация схемы БД — идемпотентно (CREATE IF NOT EXISTS) + startup cleanup.
|
||||
|
||||
Вызывается при старте приложения (в app.py, до register_blueprint).
|
||||
Безопасно для нескольких gunicorn-воркеров — PostgreSQL корректно
|
||||
обрабатывает конкурентный DDL.
|
||||
Вызывается при ПЕРВОМ обращении к БД в каждом gunicorn-воркере.
|
||||
Безопасно для конкурентного вызова из нескольких воркеров — PostgreSQL
|
||||
корректно обрабатывает одновременный CREATE IF NOT EXISTS.
|
||||
|
||||
Таблица 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 (версия приложения)
|
||||
ТАБЛИЦЫ:
|
||||
|
||||
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,
|
||||
@@ -42,17 +62,21 @@ CREATE TABLE IF NOT EXISTS runs (
|
||||
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);
|
||||
@@ -62,7 +86,7 @@ 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,
|
||||
@@ -80,16 +104,19 @@ CREATE TABLE IF NOT EXISTS scenario_runs (
|
||||
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);
|
||||
|
||||
-- Миграции для 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,
|
||||
@@ -103,16 +130,28 @@ CREATE TABLE IF NOT EXISTS scenario_definitions (
|
||||
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 при первом старте."""
|
||||
"""Импорт дефолтных сценариев из 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")
|
||||
@@ -142,7 +181,15 @@ def _seed_scenarios():
|
||||
|
||||
|
||||
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
|
||||
@@ -154,6 +201,8 @@ def init_db():
|
||||
|
||||
try:
|
||||
cur = conn.cursor()
|
||||
|
||||
# Шаг 1-2: схема + миграции
|
||||
cur.execute(SCHEMA_SQL)
|
||||
cur.execute(MIGRATION_SQL)
|
||||
cur.execute(SCENARIO_RUNS_SQL)
|
||||
@@ -161,7 +210,9 @@ def init_db():
|
||||
cur.close()
|
||||
print("[DB] Schema initialized", flush=True)
|
||||
|
||||
# Startup cleanup: зависшие RUNNING сценарии старше 1 часа → TIMEOUT
|
||||
# Шаг 3: Startup cleanup — перевести зависшие RUNNING в TIMEOUT
|
||||
# Если gunicorn упал во время выполнения сценария, статус остался RUNNING.
|
||||
# Чистим всё что старше 1 часа — операция точно не могла длиться дольше.
|
||||
try:
|
||||
cur = conn.cursor()
|
||||
cur.execute("""
|
||||
@@ -183,5 +234,5 @@ def init_db():
|
||||
finally:
|
||||
put_conn(conn)
|
||||
|
||||
# Seed сценариев — только при первом старте (INSERT ON CONFLICT DO NOTHING)
|
||||
# Шаг 4: Seed сценариев (вне основной транзакции)
|
||||
_seed_scenarios()
|
||||
|
||||
+46
-11
@@ -1,26 +1,41 @@
|
||||
"""
|
||||
Connection pool для PostgreSQL (psycopg2).
|
||||
Connection pool для PostgreSQL (psycopg2) — lazy-init в каждом gunicorn-воркере.
|
||||
|
||||
Lazy-init: пул создаётся при первом обращении к БД в каждом gunicorn-воркере.
|
||||
Почему lazy-init: gunicorn форкает воркеры, и если создать пул до fork,
|
||||
все воркеры будут делить ОДНО соединение → гонки и падения.
|
||||
Поэтому _pool создаётся при ПЕРВОМ обращении к БД в каждом воркере отдельно.
|
||||
|
||||
Переменные окружения (КАЖДАЯ отдельно):
|
||||
DB_HOST — хост
|
||||
DB_PORT — порт (5432)
|
||||
DB_NAME — имя БД
|
||||
Переменные окружения (КАЖДАЯ отдельно, не DSN-строкой):
|
||||
DB_HOST — хост PostgreSQL
|
||||
DB_PORT — порт (по умолчанию 5432)
|
||||
DB_NAME — имя базы данных
|
||||
DB_USER — пользователь
|
||||
DB_PASSWORD — пароль
|
||||
DB_SSLMODE — sslmode (require)
|
||||
DB_SSLMODE — sslmode (по умолчанию "require")
|
||||
|
||||
ThreadedConnectionPool(1, 5):
|
||||
- Минимум 1 соединение (всегда готово)
|
||||
- Максимум 5 одновременных соединений на воркер
|
||||
- При 2 воркерах gunicorn: макс. 10 соединений к БД всего
|
||||
"""
|
||||
|
||||
import os
|
||||
import psycopg2
|
||||
from psycopg2 import pool
|
||||
|
||||
# _pool — по одному на gunicorn-воркер (создаётся при первом get_conn)
|
||||
_pool = None
|
||||
# _initialized — схема уже применена в этом воркере
|
||||
_initialized = False
|
||||
|
||||
|
||||
def _dsn():
|
||||
"""Построить DSN-строку из отдельных переменных окружения.
|
||||
|
||||
DSN (Data Source Name) — это строка подключения для psycopg2:
|
||||
"host=... port=... dbname=... user=... password=... sslmode=..."
|
||||
|
||||
Каждое поле отдельной переменной — безопаснее чем один DATABASE_URL."""
|
||||
return (
|
||||
f"host={os.getenv('DB_HOST')} "
|
||||
f"port={os.getenv('DB_PORT', '5432')} "
|
||||
@@ -32,7 +47,10 @@ def _dsn():
|
||||
|
||||
|
||||
def _ensure_schema():
|
||||
"""Инициализировать схему при первом обращении к БД."""
|
||||
"""Инициализировать схему БД при первом обращении в этом воркере.
|
||||
|
||||
Вызывает init_db() который делает CREATE IF NOT EXISTS — идемпотентно.
|
||||
_initialized — глобальный флаг чтобы не дёргать init_db при каждом запросе."""
|
||||
global _initialized
|
||||
if _initialized:
|
||||
return
|
||||
@@ -41,25 +59,42 @@ def _ensure_schema():
|
||||
init_db()
|
||||
_initialized = True
|
||||
except Exception:
|
||||
pass
|
||||
pass # без БД приложение работает (без истории)
|
||||
|
||||
|
||||
def get_pool():
|
||||
"""Получить connection pool (создать при первом вызове).
|
||||
|
||||
ThreadedConnectionPool — каждый поток получает своё соединение.
|
||||
Для gunicorn с sync-воркерами (не threads) это эквивалентно SimpleConnectionPool.
|
||||
|
||||
Returns:
|
||||
ThreadedConnectionPool или None если DB_USER не задан."""
|
||||
global _pool
|
||||
if _pool is None:
|
||||
if not os.getenv("DB_USER"):
|
||||
return None
|
||||
return None # нет переменных БД — работаем без истории
|
||||
_pool = pool.ThreadedConnectionPool(1, 5, _dsn())
|
||||
_ensure_schema()
|
||||
_ensure_schema() # при первом соединении применяем схему
|
||||
return _pool
|
||||
|
||||
|
||||
def get_conn():
|
||||
"""Взять соединение из пула.
|
||||
|
||||
Вызывается перед КАЖДЫМ запросом к БД.
|
||||
Возвращает None если БД не настроена.
|
||||
|
||||
ВАЖНО: после использования ОБЯЗАТЕЛЬНО вернуть через put_conn()."""
|
||||
p = get_pool()
|
||||
return p.getconn() if p else None
|
||||
|
||||
|
||||
def put_conn(conn):
|
||||
"""Вернуть соединение в пул.
|
||||
|
||||
ВСЕГДА вызывается в finally-блоке после get_conn.
|
||||
Если не вернуть — пул исчерпается и приложение встанет."""
|
||||
p = get_pool()
|
||||
if p and conn:
|
||||
p.putconn(conn)
|
||||
|
||||
+37
-4
@@ -1,5 +1,13 @@
|
||||
"""
|
||||
Сохранение результатов тестов в БД.
|
||||
Сохранение результатов тестов в БД (таблица runs).
|
||||
|
||||
Вызывается после КАЖДОЙ операции:
|
||||
- Ручной режим: _finish_op() в api_test.py
|
||||
- Сценарный режим: run_scenario() в scenario.py (после каждого шага)
|
||||
|
||||
Сохраняет ВСЕ параметры запроса и ответа — для истории и анализа ошибок.
|
||||
Секретные поля (password, token, ...) маскируются через redact_params().
|
||||
instance_meta — полный JSON ответа GET /instances/{uid} (для post-mortem анализа).
|
||||
"""
|
||||
|
||||
import json
|
||||
@@ -11,8 +19,31 @@ def save_run(client_id, stand, user_email, svc_id, svc_name, op_name, svc_op_id,
|
||||
error_log, params, stages, app_version,
|
||||
scenario_run_id=None, step_number=None, instance_meta=None):
|
||||
"""Сохранить один запуск в таблицу runs.
|
||||
Опционально: scenario_run_id + step_number для шагов сценария.
|
||||
instance_meta — JSONB с полной информацией об инстансе."""
|
||||
|
||||
Эта функция — единственная точка записи в runs.
|
||||
Вызывается и из ручного режима, и из сценариев.
|
||||
|
||||
Args:
|
||||
client_id: str — ClientID пользователя
|
||||
stand: str — "dev"/"test"
|
||||
user_email: str — email из JWT
|
||||
svc_id: int — ID сервиса
|
||||
svc_name: str — имя сервиса (из ответа poll)
|
||||
op_name: str — "create"/"modify"/"delete"/"suspend"/"resume"/"redeploy"
|
||||
svc_op_id: int — числовой ID операции
|
||||
op_uid: str — UUID операции
|
||||
instance_uid: str — UUID инстанса
|
||||
display_name: str — displayName инстанса
|
||||
status: str — "OK"/"FAIL"/"TIMEOUT"/"RUNNING"
|
||||
duration_sec: float — длительность операции
|
||||
error_log: str — текст ошибки (если FAIL)
|
||||
params: dict — параметры операции (уже маскированные)
|
||||
stages: list — этапы выполнения
|
||||
app_version: str — версия приложения
|
||||
scenario_run_id: int|None — ID запуска сценария
|
||||
step_number: int|None — номер шага в сценарии
|
||||
instance_meta: dict|None — полная информация об инстансе (JSONB)
|
||||
"""
|
||||
conn = get_conn()
|
||||
if not conn:
|
||||
print("[DB] save_run SKIP: no connection", flush=True)
|
||||
@@ -30,10 +61,12 @@ def save_run(client_id, stand, user_email, svc_id, svc_name, op_name, svc_op_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 и stages — JSONB: сериализуем в JSON-строку
|
||||
json.dumps(params) if params else None,
|
||||
json.dumps(stages) if stages else None,
|
||||
app_version,
|
||||
scenario_run_id, step_number,
|
||||
# instance_meta — тоже JSONB
|
||||
json.dumps(instance_meta) if instance_meta else None,
|
||||
))
|
||||
conn.commit()
|
||||
@@ -43,4 +76,4 @@ def save_run(client_id, stand, user_email, svc_id, svc_name, op_name, svc_op_id,
|
||||
print(f"[DB] save_run error: {e}", flush=True)
|
||||
conn.rollback()
|
||||
finally:
|
||||
put_conn(conn)
|
||||
put_conn(conn) # ВСЕГДА возвращаем соединение в пул
|
||||
|
||||
Reference in New Issue
Block a user