v1.2.20: remaining Codex fixes — _op_results lock, advisory lock for scenarios, _ensure_schema logging, stale async generation token, validate-cfs explicit JSONDecodeError

This commit is contained in:
2026-07-31 18:07:24 +04:00
parent 58b823a260
commit 5065019ffd
7 changed files with 222 additions and 172 deletions
+1 -1
View File
@@ -32,7 +32,7 @@ from routes.api_scenario_defs import bp_defs as api_scenario_defs_bp
# Версия — показывается в топбаре UI. Меняется при КАЖДОМ изменении кода.
# Нужна для фильтрации истории (пользователь видит только записи своей версии).
VERSION = "1.2.19"
VERSION = "1.2.20"
# Flask-приложение с Jinja2-шаблонами из папки templates/
app = Flask(__name__, template_folder="templates", static_folder="static")
+5 -2
View File
@@ -58,8 +58,11 @@ def _ensure_schema():
from db.init_db import init_db
init_db()
_initialized = True
except Exception:
pass # без БД приложение работает (без истории)
except Exception as e:
print(f"[DB] Schema init FAILED: {e}", flush=True)
import traceback
traceback.print_exc()
# Приложение продолжает работу БЕЗ БД (без истории/сценариев)
def get_pool():
+41 -12
View File
@@ -209,27 +209,56 @@ def delete_definition(def_id, client_id, stand):
def lock_check(client_id, stand):
"""Проверить что нет активного RUNNING-запуска сценария.
Используется перед запуском нового сценария — предотвращает
одновременный запуск двух сценариев одним пользователем.
Использует pg_try_advisory_lock для АТОМАРНОЙ проверки-и-захвата.
Это устраняет race condition между lock_check и INSERT scenario_runs:
если advisory lock взят — никто другой не сможет запустить сценарий
для этого же client_id+stand, пока мы его не отпустим.
Returns:
True — можно запускать (нет RUNNING)
False — нельзя (уже есть RUNNING) → 409 Conflict"""
True — можно запускать (lock взят)
False — нельзя (другой процесс держит lock) → 409 Conflict"""
conn = get_conn()
if not conn:
return True # без БД — разрешаем (fallback)
try:
cur = conn.cursor()
cur.execute("""
SELECT id FROM scenario_runs
WHERE client_id = %s AND stand = %s AND status = 'RUNNING'
LIMIT 1
""", (client_id, stand))
row = cur.fetchone()
# hashtext даёт стабильный int из строки — одинаковый во всех сессиях
lock_key = f"scenario:{client_id}:{stand}"
cur.execute("SELECT pg_try_advisory_lock(hashtext(%s))", (lock_key,))
acquired = cur.fetchone()[0]
cur.close()
return row is None # None = нет RUNNING = можно запускать
if not acquired:
put_conn(conn)
return False # lock уже взят другим процессом
# НЕ возвращаем conn в пул! Держим до unlock.
# conn будет передан в сценарий и освобождён после завершения.
return True
except Exception as e:
print(f"[DEFS] lock_check error: {e}", flush=True)
return True
put_conn(conn)
return True # fallback
def unlock_scenario(client_id, stand, conn=None):
"""Освободить advisory lock после завершения сценария.
Вызывается из run_scenario (scenario.py) после завершения ВСЕХ шагов.
conn — то же соединение, на котором был взят lock (если есть)."""
if not conn:
conn = get_conn()
if not conn:
return
try:
lock_key = f"scenario:{client_id}:{stand}"
cur = conn.cursor()
cur.execute("SELECT pg_advisory_unlock(hashtext(%s))", (lock_key,))
cur.close()
conn.commit()
except Exception as e:
print(f"[DEFS] unlock_scenario error: {e}", flush=True)
try:
conn.rollback()
except Exception:
pass
finally:
put_conn(conn)
+5
View File
@@ -32,6 +32,7 @@ from operations.executor import execute_operation
from operations.poll import poll_until_done
from db.pool import get_conn, put_conn
from db.save_run import save_run
from db.scenario_defs import unlock_scenario
# Все инстансы созданные сценарием имеют такой префикс в displayName
AUTOTEST_PREFIX = "autotest-scenario-"
@@ -174,6 +175,7 @@ def run_scenario(client, steps, client_id, stand, user_email, app_version, scena
bindings = {} # output_name → instance_uid (новый формат, именованные ссылки)
instance_map = {} # service_id → instance_uid (fallback, старый формат)
try:
for i, step in enumerate(steps):
step_num = i + 1 # шаги нумеруются с 1 (для пользователя)
svc_id = step.get("service_id")
@@ -306,3 +308,6 @@ def run_scenario(client, steps, client_id, stand, user_email, app_version, scena
# Все шаги пройдены успешно
_save_scenario_run(scenario_run_id, "OK", total, round(time.time() - t0, 1))
finally:
# Освободить advisory lock ВСЕГДА (даже при исключении)
unlock_scenario(client_id, stand)
+6 -8
View File
@@ -245,13 +245,11 @@ def send_params_terraform(client, op_uid, params):
# Шаг 7: validate-cfs — финальная проверка
# Если параметры невалидны — API вернёт ошибку → исключение
# Если ответ пустой или не-JSON — это ОК (значит валидация прошла)
# Если ответ 200 с пустым телом — это ОК (значит валидация прошла).
# HttpClient.get() для пустого тела возвращает {} без ошибок.
# Но если тело не-JSON — будет JSONDecodeError, это тоже ОК для validate-cfs.
try:
client.get(f"/instanceOperations/{op_uid}/validate-cfs")
except Exception as e:
# "Expecting value" / JSONDecodeError — это норма для validate-cfs
# (API возвращает 200 с пустым телом при успехе)
if "Expecting value" in str(e) or "JSON" in str(type(e).__name__):
pass
else:
raise # реальная ошибка — пробрасываем выше
except json.JSONDecodeError:
pass # пустой/не-JSON ответ — норма для validate-cfs (успех)
# HTTPError (4xx/5xx) пробрасывается выше — это реальная ошибка валидации
+15 -8
View File
@@ -24,6 +24,7 @@
"""
from flask import Blueprint, current_app, jsonify, request
import threading
import uuid
import os
import json
@@ -261,23 +262,27 @@ def api_test():
# Результаты фоновых операций: opUid → {status, error, stages, duration, _ts}
# _ts — timestamp добавления, для TTL-очистки (макс. 500 записей или старше 1 часа)
# ЗАЩИЩЕНО _op_results_lock — несколько потоков _finish_op + main thread api_test_status
_op_results = {}
_op_results_lock = threading.Lock()
_MAX_OP_RESULTS = 500
def _cleanup_op_results():
"""Удалить старые записи: старше 1 часа или сверх лимита."""
"""Удалить старые записи: старше 1 часа или сверх лимита.
Потокобезопасно — под _op_results_lock."""
import time
now = time.time()
# Удалить старше 1 часа
with _op_results_lock:
# Удалить старше 1 часа (pop с default — безопасно при конкурентном доступе)
stale = [k for k, v in _op_results.items() if now - v.get("_ts", 0) > 3600]
for k in stale:
del _op_results[k]
_op_results.pop(k, None)
# Если всё ещё много — удалить самые старые
if len(_op_results) > _MAX_OP_RESULTS:
sorted_keys = sorted(_op_results.keys(), key=lambda k: _op_results[k].get("_ts", 0))
for k in sorted_keys[:len(_op_results) - _MAX_OP_RESULTS]:
del _op_results[k]
_op_results.pop(k, None)
def _get_instance_display_name(client, instance_uid):
@@ -291,7 +296,7 @@ def _get_instance_display_name(client, instance_uid):
def _finish_op(client, op_uid, instance_uid, svc_id, display_name, op_name, svc_op_id, is_create, client_id, stand, is_delete=False, params=None, user_email="", app_version=""):
"""Фоном ждать dtFinish и сохранить результат."""
"""Фоном ждать dtFinish и сохранить результат (потокобезопасно для _op_results)."""
import time
_cleanup_op_results()
t0 = time.time()
@@ -299,7 +304,8 @@ def _finish_op(client, op_uid, instance_uid, svc_id, display_name, op_name, svc_
# Поллинг через общий модуль
poll_result = poll_until_done(client, op_uid)
# Обновить _op_results для UI
# Обновить _op_results для UI — под локом
with _op_results_lock:
_op_results[op_uid] = {
"status": poll_result["status"],
"displayName": display_name,
@@ -346,8 +352,9 @@ def _finish_op(client, op_uid, instance_uid, svc_id, display_name, op_name, svc_
@bp.route("/api/test/status/<op_uid>")
def api_test_status(op_uid):
"""Получить текущий статус операции (поллинг с UI)."""
# сначала проверяем фоновый трекер
"""Получить текущий статус операции (поллинг с UI) — потокобезопасно."""
# сначала проверяем фоновый трекер — под локом
with _op_results_lock:
if op_uid in _op_results:
return jsonify(_op_results[op_uid])
# иначе спрашиваем API напрямую
+9 -1
View File
@@ -208,6 +208,9 @@ async function loadStepParams(idx) {
const step = st.steps[idx];
if (!step.operation) return;
// Запомнить поколение рендера на момент старта async-загрузки
const renderGen = st._renderGen || 0;
const svcOpId = await resolveStepSvcOpId(idx);
if (!svcOpId) return;
step._svcOpId = svcOpId; // сохраняем для обратного маппинга при сохранении
@@ -233,6 +236,8 @@ async function loadStepParams(idx) {
// Рендер: клонируем def и подставляем сохранённые значения в defaultValue
const container = document.getElementById('step-' + idx + '-params');
if (!container) return;
// Проверка: не перезаписывать DOM если уже был новый renderEditor()
if (st._renderGen !== renderGen) return;
let html = '';
defs.forEach(p => {
const clone = Object.assign({}, p); // shallow copy
@@ -259,7 +264,10 @@ async function loadStepParams(idx) {
function renderEditor() {
snapshotAllParams(); // сохранить несохранённые правки перед перерисовкой
const st = scenarioEditorState;
// Инкремент поколения — защита от stale async (loadStepParams может вернуться
// уже после следующего renderEditor и перезаписать новый DOM старыми данными)
st._renderGen = (st._renderGen || 0) + 1;
const currentGen = st._renderGen;
let html = '<div style="padding:20px;max-height:90vh;overflow-y:auto;">';
// Заголовок