diff --git a/site/api/auth.py b/site/api/auth.py index a7bb974..4898b3d 100644 --- a/site/api/auth.py +++ b/site/api/auth.py @@ -3,15 +3,15 @@ Раньше _client(), _client_id(), _stand(), _token_info() были продублированы в main.py и api_test.py с идентичным или почти идентичным кодом. -Теперь всё здесь — один источник правды. +Теперь всё здесь — один источник правды для всех роутов. -Функции: - get_token() — токен из cookie или env +Функции (все без аргументов — берут данные из Flask request/current_app): + get_token() — токен: cookie → env-переменная get_client() — HttpClient с автоопределением стенда - get_client_id() — ClientID из JWT (base64, без проверки подписи) - get_stand() — "dev"/"test" по токену - get_token_info() — {email, company, client_id} из JWT - get_token_masked() — маскированный токен (abc...xyz) + get_client_id() — ClientID из JWT (base64url, без проверки подписи) + get_stand() — "dev" / "test" по токену + get_token_info() — {email, company, client_id} из JWT для UI + get_token_masked() — маскированный токен (abc...xyz) для placeholder """ import base64 @@ -22,13 +22,19 @@ from api.http_client import HttpClient, detect_endpoint, stand_name def get_token(): - """Токен: сначала из cookie, потом из env-переменной.""" + """Получить активный токен: сначала из cookie пользователя, потом из env. + + Приоритет: + 1. cookie "token" — пользователь ввёл свой токен в форме + 2. NUBES_API_TOKEN из env — сервисный токен (для автоматических тестов) + + Пользовательский токен приоритетнее — он переопределяет сервисный.""" return request.cookies.get("token") or current_app.config["NUBES_API_TOKEN"] def get_client(): - """HttpClient с автоопределением стенда по токену. - + """HttpClient с автоопределением стенда по активному токену. + Использует detect_endpoint() — пробует dev→test стенды. Если автоопределение не сработало — fallback на NUBES_API_ENDPOINT из конфига.""" token = get_token() @@ -37,27 +43,45 @@ def get_client(): def get_client_id(): - """Извлечение ClientID из payload JWT-токена (base64url, без проверки подписи).""" + """Извлечь ClientID из payload JWT-токена (base64url, без проверки подписи). + + JWT состоит из трёх частей: header.payload.signature, разделённых точкой. + Нам нужен ТОЛЬКО payload — он в base64url (не base64!). + + ClientID используется для: + - Изоляции данных в БД (runs.client_id) + - Изоляции трекера инстансов (/tmp/instances-{clientId}-{stand}.json) + - Отображения в UI + + Безопасность: мы НЕ проверяем подпись — это не нужно. + Токен уже проверен Nubes API (detect_endpoint делает реальный запрос).""" token = get_token() try: - parts = token.split(".") # header.payload.signature + parts = token.split(".") # [header, payload, signature] if len(parts) >= 2: - payload = base64.urlsafe_b64decode(parts[1] + "==") # padding + # base64url → добавляем padding ("==") на случай если длина не кратна 4 + payload = base64.urlsafe_b64decode(parts[1] + "==") return json.loads(payload).get("ClientID", "") except Exception: - pass + pass # битый токен — не критично, вернём пустую строку return "" def get_stand(): - """dev/test — по токену (detect_endpoint → stand_name).""" + """Определить стенд (dev/test) по активному токену. + + detect_endpoint → stand_name. Если автоопределение не сработало — + fallback на NUBES_API_ENDPOINT из конфига.""" token = get_token() endpoint = detect_endpoint(token) or current_app.config["NUBES_API_ENDPOINT"] return stand_name(endpoint) def get_token_info(): - """{email, company, client_id} из JWT — для отображения в топбаре UI.""" + """Извлечь {email, company, client_id} из JWT — для отображения в топбаре UI. + + Возвращает dict с ключами: email, company, client_id. + Если JWT невалиден — возвращает пустой {}.""" token = get_token() try: parts = token.split(".") @@ -75,8 +99,12 @@ def get_token_info(): def get_token_masked(): - """Маскированный токен для placeholder: abc...xyz.""" + """Маскированный env-токен для placeholder в форме: abc...xyz. + + Используется ТОЛЬКО env-токен (не пользовательский!). + Если токен короче 8 символов — возвращает пустую строку.""" token = current_app.config["NUBES_API_TOKEN"] if not token or len(token) < 8: return "" + # Первые 4 символа + звёздочки + последние 4 символа return token[:4] + "*" * (len(token) - 8) + token[-4:] diff --git a/site/api/http_client.py b/site/api/http_client.py index 59be952..f63308d 100644 --- a/site/api/http_client.py +++ b/site/api/http_client.py @@ -1,20 +1,27 @@ """ -HTTP-клиент для Nubes API + автоопределение стенда. +HTTP-клиент для Nubes API — тонкая обёртка над requests.Session. -Класс HttpClient — тонкая обёртка над requests.Session: - - Добавляет заголовки: Authorization Bearer, User-Agent (DDoS-Guard) - - GET: raise_for_status → .json() - - POST: проверка r.ok, извлечение Location-заголовка +HttpClient: + - Добавляет обязательные заголовки: + Authorization: Bearer — аутентификация + User-Agent: Mozilla/5.0 — DDoS-Guard блокирует python-requests по умолчанию + - GET: raise_for_status → .json() — автоматически проверяет HTTP-статус + - POST: проверка r.ok, извлечение Location-заголовка + UUID + - raw_delete: DELETE без авторизации (для CMDB) -Функции автостенда: - - detect_endpoint(token) — пробует dev→test стенды, возвращает URL - - stand_name(endpoint) — "dev" / "test" по URL - - create_client(token, fallback) — HttpClient + endpoint (одним вызовом) +Функции автоопределения стенда: + - detect_endpoint(token) — пробует dev→test стенды по токену + - stand_name(endpoint) — "dev" / "test" по URL + - create_client(token) — HttpClient + endpoint одним вызовом + +Зачем автоопределение: пользователь вводит токен, мы не знаем dev это или test. +Пробуем оба стенда — какой ответит с results != None, тот и рабочий. """ import requests -# Список стендов для автоопределения. Порядок важен: dev первый (быстрее). +# Список стендов для автоопределения. +# Порядок ВАЖЕН: dev первый — он быстрее (меньше нагрузка), test — резервный. STANDS = [ "https://lk-api-gateway-dev.ngcloud.ru/api/v1/svc", "https://lk-api-gateway-test.ngcloud.ru/api/v1/svc", @@ -23,10 +30,16 @@ STANDS = [ def detect_endpoint(token): """Пробуем токен против dev и test стендов, возвращаем рабочий URL. - - Делает GET /instances?pageSize=1 на каждый стенд. - Если results != None — стенд рабочий, возвращаем его URL. - Если ни один не подошёл — возвращаем None.""" + + Алгоритм: + 1. Для каждого стенда создаём HttpClient + делаем GET /instances?pageSize=1. + 2. Если results != None — стенд рабочий, токен валиден → возвращаем URL. + 3. Если Exception (401/403/таймаут) — пробуем следующий. + 4. Ни один не подошёл → None. + + Почему results != None а не просто HTTP 200: + API может вернуть 200 с пустым списком (нет инстансов) — это норм. + results=None — признак что ответ не соответствует ожидаемой структуре.""" for ep in STANDS: try: c = HttpClient(ep, token) @@ -39,10 +52,10 @@ def detect_endpoint(token): def stand_name(endpoint): - """dev/test по URL стенда. - + """Определить имя стенда по URL. + Ищет подстроку 'dev' или 'test' в URL. - Если не найдено — возвращает '?'.""" + Если не найдено — '?' (неизвестный стенд).""" for name in ("dev", "test"): if name in (endpoint or ""): return name @@ -50,10 +63,11 @@ def stand_name(endpoint): def create_client(token, fallback_endpoint=None): - """HttpClient с автоопределением стенда по токену. - - Сначала detect_endpoint(token), если не сработало — fallback_endpoint. - Возвращает кортеж (HttpClient, endpoint_url) или None если стенд не определён.""" + """HttpClient с автоопределением стенда. + + Сначала detect_endpoint(token) — dev→test. + Если не сработало — fallback_endpoint (из конфига). + Возвращает (HttpClient, endpoint_url) или None.""" ep = detect_endpoint(token) or fallback_endpoint if not ep: return None @@ -62,43 +76,66 @@ def create_client(token, fallback_endpoint=None): class HttpClient: """HTTP-клиент для Nubes REST API. - - Использует requests.Session для keep-alive соединений. - Добавляет обязательные заголовки: - - Authorization: Bearer (аутентификация) - - User-Agent: Mozilla/5.0 (DDoS-Guard требует НЕ python-requests)""" - + + Использует requests.Session для: + - Keep-alive (переиспользование TCP+TLS между запросами) + - Единые заголовки для всех запросов + + Обязательные заголовки: + - Authorization: Bearer — без этого API вернёт 401 + - User-Agent: Mozilla/5.0 — DDoS-Guard блокирует "python-requests/2.x" + + Таймауты: GET=10с (лёгкие), POST=30с (создание ресурсов дольше).""" + def __init__(self, endpoint, token): - # Убираем trailing slash чтобы потом добавлять "/path" + # Убираем trailing slash (".../svc/" → ".../svc") + # чтобы path добавлялся как "/path", а не "path" self._endpoint = endpoint.rstrip("/") - # Session — переиспользует TCP-соединения между запросами + # Session переиспользует TCP + TLS handshake между запросами. + # Без Session каждый запрос делал бы новый connect → медленно. self._session = requests.Session() self._session.headers.update({ "Authorization": f"Bearer {token}", - "User-Agent": "Mozilla/5.0", + "User-Agent": "Mozilla/5.0", # DDoS-Guard: не python-requests }) def get(self, path, **kwargs): - """GET-запрос. Таймаут по умолчанию 10 секунд. - - Вызывает raise_for_status() — при 4xx/5xx выбрасывает HTTPError. - Возвращает распарсенный JSON (dict/list).""" + """GET-запрос к Nubes API. + + Особенности: + - Таймаут 10 секунд по умолчанию. + - raise_for_status() — 4xx/5xx → HTTPError (ловим выше). + - Пустой ответ (нет body) → {} (норма для validate-cfs). + + Args: + path: str — путь относительно endpoint (напр. "/instances"). + **kwargs — params, timeout, headers и т.д. + + Returns: + dict/list — распарсенный JSON, или {} если ответ пустой.""" kwargs.setdefault("timeout", 10) r = self._session.get(f"{self._endpoint}{path}", **kwargs) r.raise_for_status() if not r.text or not r.text.strip(): - return {} # пустой ответ (напр. validate-cfs успех) + return {} # пустой ответ — норма для некоторых эндпоинтов return r.json() def post(self, path, data=None, **kwargs): - """POST-запрос. Таймаут по умолчанию 30 секунд. - - Отправляет data как JSON (json=...). - При HTTP-ошибке выбрасывает Exception с кодом и телом ответа. - Возвращает dict: - - Распарсенный JSON (если ответ — валидный JSON-объект) - - + ключ "_location" со значением заголовка Location (если есть) - Location нужен для получения UID созданного ресурса (instanceUid, opUid).""" + """POST-запрос к Nubes API. + + Особенности: + - Таймаут 30 секунд по умолчанию. + - data → json=... — requests сам ставит Content-Type: application/json. + - Извлекает Location-заголовок → UUID → кладёт в instanceUid/instanceOperationUid. + Это нужно чтобы не полагаться только на тело ответа (которое может быть пустым). + - При HTTP-ошибке: Exception с кодом и телом ответа. + + Args: + path: str — путь (напр. "/instances"). + data: dict|None — тело запроса (сериализуется в JSON). + + Returns: + dict с полями ответа + _status + _location (+ UUID если найден).""" kwargs.setdefault("timeout", 30) url = f"{self._endpoint}{path}" # json=... — requests сам сериализует и ставит Content-Type: application/json @@ -112,15 +149,17 @@ class HttpClient: result = parsed result["_status"] = r.status_code except Exception: - pass # тело не JSON или не dict — ок, Location всё равно извлечём + pass # тело не JSON — ок, UUID извлечём из Location + # Location-заголовок: "./UUID" или "/api/v1/svc/.../UUID" loc = r.headers.get("Location", "") if loc: - # Location бывает вида "./UUID" или "/api/v1/svc/instanceOperations/UUID" result["_location"] = loc - # Извлечь UUID и положить в правильное поле ответа + # Извлекаем UUID — последний сегмент после split("/") parts = loc.rstrip("/").split("/") uid = parts[-1] + # Проверка: не ".", длина >= 32 (UUID = 36 символов с дефисами) if uid and uid != "." and len(uid) >= 32: + # Кладём UUID в правильное поле в зависимости от эндпоинта if "/instanceOperations" in path: result["instanceOperationUid"] = uid elif "/instances" in path: @@ -128,9 +167,17 @@ class HttpClient: return result def raw_delete(self, url): - """DELETE-запрос к произвольному URL (CMDB API). - Использует отдельную сессию БЕЗ auth-заголовков. - Возвращает кортеж (ok: bool, status_code: int).""" + """DELETE-запрос к произвольному URL — БЕЗ авторизации. + + Используется для CMDB API (cmdb-api.deck.nubes.ru) — + жёсткое удаление недосозданных/зависших инстансов. + CMDB не требует Bearer-токена. + + Args: + url: str — полный URL (не path, другой хост!) + + Returns: + (ok: bool, status_code: int|str).""" try: r = requests.delete(url, timeout=10) return r.ok, r.status_code diff --git a/site/api/utils.py b/site/api/utils.py index 43ac156..c6b2608 100644 --- a/site/api/utils.py +++ b/site/api/utils.py @@ -1,45 +1,103 @@ """ Общие утилиты: извлечение UUID из ответов API Nubes. -Используется: api_test.py, scenario.py, executor.py. -Заменяет разрозненные реализации _find_uid / _uid_from_location. +Проблема: Nubes API возвращает UUID инстанса/операции в разных форматах +в зависимости от эндпоинта и фазы жизненного цикла: + - POST /instances → {"instanceUid": "..."} (прямой ключ) + - POST /instances → Location: "./uuid" (заголовок) + - POST /instanceOperations → {"instanceOperationUid": "..."} (другой ключ) + - GET /instanceOperations → {"instanceOperation": {"instanceOperationUid": "..."}} (вложенный) + +Эти две функции — единый способ достать UUID из ЛЮБОГО ответа. +Раньше каждая точка вызова имела свою реализацию _find_uid / _uid_from_location. +Теперь всё здесь — один источник правды. + +Используется: executor.py, api_test.py, scenario.py. """ def find_uid(resp): - """Извлечь instanceUid или instanceOperationUid из ответа API. + """Извлечь instanceUid или instanceOperationUid из dict-ответа API. - Ищет в нескольких местах: - 1. На верхнем уровне: instanceOperationUid → instanceUid → uid → Uid - 2. Во вложенных dict-значениях: {key: {instanceOperationUid: ..., instanceUid: ...}} + Стратегия поиска (в порядке приоритета): + 1. Верхний уровень — прямые ключи (стиль api_test.py): + instanceOperationUid → instanceUid → uid → Uid + 2. Вложенные dict-значения (стиль scenario.py): + Для каждого значения-словаря проверяем instanceOperationUid → instanceUid + + Почему такой порядок: + - instanceOperationUid приоритетнее — это ответ на POST /instanceOperations + - instanceUid — ответ на POST /instances или GET /instances + - uid / Uid — legacy-форматы, почти не встречаются + + Args: + resp: dict — распарсенный JSON-ответ API (или что угодно). + + Returns: + str или None — UUID (36 символов) или None если не нашли. """ + # Защита: если передали не dict (например, list или str) — сразу None if not isinstance(resp, dict): return None - # Верхний уровень — прямые ключи (api_test.py-стиль) + + # ── Уровень 1: прямые ключи на верхнем уровне ── + # Пример: {"instanceUid": "abc-123", ...} + # Проверяем по порядку: сначала специфичные, потом общие for key in ("instanceOperationUid", "instanceUid", "uid", "Uid"): v = resp.get(key) + # isinstance(str) — отсекаем числа, None, пустые строки if isinstance(v, str) and v: return v - # Вложенные dict-значения (scenario.py-стиль) + + # ── Уровень 2: вложенные dict-значения ── + # Пример: {"instanceOperation": {"instanceOperationUid": "abc-123"}} + # Перебираем ВСЕ значения resp, ищем вложенные словари for v in resp.values(): if isinstance(v, dict): + # Внутри вложенного dict — те же ключи что на верхнем уровне uid = v.get("instanceOperationUid") or v.get("instanceUid") if uid: return uid + + # Ничего не нашли return None def uid_from_location(loc): - """Извлечь UUID из Location-заголовка. + """Извлечь UUID из Location-заголовка HTTP-ответа. - Location бывает: "./UUID", "/api/v1/svc/instances/UUID", "/api/v1/svc/instanceOperations/UUID" - Берём последний сегмент после split("/"). + Location — это HTTP-заголовок, который Nubes API возвращает при создании + ресурса (201 Created). Форматы: + - "./uuid" (относительный) + - "/api/v1/svc/instances/uuid" (абсолютный путь) + - "/api/v1/svc/instanceOperations/uuid" (полный путь) + + Алгоритм: + 1. Убираем trailing slash. + 2. Split по "/" — берём ПОСЛЕДНИЙ сегмент. + 3. Проверяем что это похоже на UUID (>= 32 символа) и не ".". + + Args: + loc: str или None — значение заголовка Location. + + Returns: + str или None — UUID или None если не нашли. """ + # Пустой Location → нечего извлекать if not loc: return None + + # Убираем trailing slash (напр. "./uuid/" → "./uuid") + # Split по "/" → [".", "uuid"] или ["", "api", "v1", "svc", "instances", "uuid"] parts = str(loc).rstrip("/").split("/") + + # Последний сегмент — это UUID uid = parts[-1] - # UUID — 36 символов с 4 дефисами (уже проверено в http_client.py) + + # Проверки: + # uid != "." — Location "./" даст пустой последний сегмент + # len >= 32 — UUID всегда 36 символов (8-4-4-4-12), но на всякий случай >= 32 if uid and uid != "." and len(uid) >= 32: return uid + return None diff --git a/site/db/init_db.py b/site/db/init_db.py index 6426df0..9bc7345 100644 --- a/site/db/init_db.py +++ b/site/db/init_db.py @@ -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() diff --git a/site/db/pool.py b/site/db/pool.py index 947543f..b7f4c00 100644 --- a/site/db/pool.py +++ b/site/db/pool.py @@ -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) diff --git a/site/db/save_run.py b/site/db/save_run.py index c4aee67..5bbe2bb 100644 --- a/site/db/save_run.py +++ b/site/db/save_run.py @@ -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) # ВСЕГДА возвращаем соединение в пул diff --git a/site/operations/executor.py b/site/operations/executor.py index 938799f..4fcdc96 100644 --- a/site/operations/executor.py +++ b/site/operations/executor.py @@ -1,9 +1,20 @@ """ -Единый executor операции Nubes. +Единый executor операции Nubes — ЕДИНСТВЕННАЯ точка запуска операций. -Делает всё до /run включительно. НЕ поллит. -Вызывает tracker_add для create сразу после получения instance_uid (защита от сирот). -Используется api_test.py (ручной) и scenario.py (сценарный). +Выполняет ВСЕ шаги от начала до /run (включительно): + 1. CREATE: POST /instances → получает instanceUid + 2. CREATE + NON-CREATE: POST /instanceOperations → получает opUid + 3. Параметры: send_params_terraform → нормализация + refSvc-резолв + отправка + 4. POST /instanceOperations/{opUid}/run → запуск + +НЕ делает поллинг — поллинг вынесен в poll.py (используется отдельно). + +Критический инвариант: tracker_add вызывается СРАЗУ после получения instanceUid +(до params и run), чтобы даже при падении gunicorn инстанс не стал сиротой. + +Используется: + - api_test.py → ручной режим (фронтенд → POST /api/test → execute_operation) + - scenario.py → сценарный режим (каждый шаг → execute_operation) """ from api.utils import find_uid, uid_from_location @@ -14,75 +25,145 @@ from operations.tracker import add as tracker_add def execute_operation(client, service_id, operation, instance_uid, params, svc_op_id=None, display_name=None, descr=None, client_id="", stand=""): - """Запустить операцию и вернуть результат. + """Запустить операцию и вернуть результат (без поллинга). + + Эта функция — «единый шлюз»: и ручной тест, и сценарий идут через неё. + Это гарантирует что поведение create/modify/delete идентично в обоих режимах. Args: - client: HttpClient - service_id: int — ID сервиса - operation: str — create/delete/modify/suspend/resume/redeploy - instance_uid: str|None — UUID существующего инстанса (None для create) - params: dict — {numeric_param_id: value} - svc_op_id: int|None — svcOperationId (игнорируется для create) - display_name: str|None — displayName (только для create) - descr: str|None — описание инстанса - client_id: str — для tracker_add - stand: str — для tracker_add + client: HttpClient (уже с токеном и endpoint) + service_id: int — ID сервиса Nubes (напр. 1 для Болванки) + operation: str — "create" | "modify" | "delete" | "suspend" | "resume" | "redeploy" + instance_uid: str|None — UUID существующего инстанса. + None для create (инстанса ещё нет). + params: dict — {numeric_param_id: string_value} + Пример: {"198": "0", "199": "test-value"} + svc_op_id: int|None — svcOperationId (числовой ID операции в Nubes). + Для create игнорируется (там свой ID из шаблона сервиса). + display_name: str|None — displayName для create. + Если None → генерируется "autotest-{service_id}". + descr: str|None — описание инстанса (только для create). + Если None → "created by autotest". + client_id: str — ID пользователя (для tracker_add). + stand: str — "dev" / "test" (для tracker_add). Returns: - dict {ok, error, failed_step, instance_uid, op_uid, display_name} - ok=True при успехе, ok=False при ошибке с failed_step. + dict с ВСЕГДА одинаковыми ключами: + {ok, error, failed_step, instance_uid, op_uid, display_name} + + ok=True → успех, можно поллить + ok=False → ошибка на шаге failed_step: + "instances" — не смогли создать инстанс (POST /instances) + "instanceOperations"— не смогли создать операцию (POST /instanceOperations) + "params" — ошибка в send_params_terraform + "run" — не смогли запустить (/run) """ + # ── Определяем режим: create или нет ── + # Все не-create операции требуют существующий instance_uid is_create = (operation == "create") - # --- Шаг 1: POST /instances (только create) --- + # ═══════════════════════════════════════════════════════ + # Шаг 1: POST /instances — только для create + # ═══════════════════════════════════════════════════════ + # Создаём инстанс в Nubes, получаем его UUID. + # Это аналог нажатия «Создать» в личном кабинете. if is_create: + # display_name — обязателен для create. + # Если не передан — генерируем из service_id. if not display_name: display_name = f"autotest-{service_id}" + + # descr — описание, видно в личном кабинете. if not descr: descr = "created by autotest" + + # payload для POST /instances: + # serviceId — какой сервис создаём + # displayName — имя инстанса (уникально в рамках пользователя) + # descr — текстовое описание payload = {"serviceId": service_id, "displayName": display_name, "descr": descr} + try: resp = client.post("/instances", payload) except Exception as e: + # Сетевая ошибка или HTTP-ошибка → сразу FAIL return {"ok": False, "error": str(e), "failed_step": "instances", "instance_uid": None, "op_uid": None, "display_name": display_name} + + # Извлекаем instanceUid из ответа API. + # Nubes API может вернуть UUID в разных местах ответа: + # 1. resp["instanceUid"] — прямой ключ + # 2. resp["instance"]["instanceUid"] — вложенный (через find_uid) + # 3. Location-заголовок — "./uuid" (через uid_from_location) instance_uid = resp.get("instanceUid") or find_uid(resp) or uid_from_location(resp.get("_location", "")) + + # Если ни один метод не дал UUID — это баг API, не можем продолжать if not instance_uid: return {"ok": False, "error": "No instanceUid in response", "failed_step": "instances", "instance_uid": None, "op_uid": None, "display_name": display_name} - # tracker_add СРАЗУ после instance_uid, до params/run — защита от сирот + # ⚠️ КРИТИЧЕСКИ: tracker_add СРАЗУ после получения instanceUid. + # Зачем: если gunicorn упадёт между create и run, инстанс останется в Nubes + # но не будет виден в UI (GET /instances иногда задерживает новый инстанс). + # Трекер — краткосрочный fallback, пока облако не подхватит. + # try/except: не ронять операцию если /tmp/ переполнен или нет прав. try: tracker_add(client_id, stand, instance_uid, service_id, display_name) except Exception: - pass # не ронять операцию из-за трекера + pass # silently ignore — трекер не критичен для операции - # --- Шаг 2: POST /instanceOperations --- + # ═══════════════════════════════════════════════════════ + # Шаг 2: POST /instanceOperations — создаём операцию + # ═══════════════════════════════════════════════════════ + # Для create: {instanceUid, operation} — svcOperationId не нужен + # Для не-create: {instanceUid, svcOperationId, operation} — нужен ID операции + # + # Разница: для create система сама знает какой svcOperationId использовать + # (он единственный для create у каждого сервиса). + # Для modify/delete/suspend/resume — их много, нужно указать конкретный. if is_create: op_payload = {"instanceUid": instance_uid, "operation": operation} else: op_payload = {"instanceUid": instance_uid, "svcOperationId": svc_op_id, "operation": operation} + try: op_resp = client.post("/instanceOperations", op_payload) except Exception as e: return {"ok": False, "error": str(e), "failed_step": "instanceOperations", "instance_uid": instance_uid, "op_uid": None, "display_name": display_name} + + # Извлекаем opUid — UUID операции, нужен для поллинга и /run op_uid = op_resp.get("instanceOperationUid") or find_uid(op_resp) or uid_from_location(op_resp.get("_location", "")) + if not op_uid: return {"ok": False, "error": "No opUid in response", "failed_step": "instanceOperations", "instance_uid": instance_uid, "op_uid": None, "display_name": display_name} - # --- Шаг 3-6: Параметры + run --- + # ═══════════════════════════════════════════════════════ + # Шаг 3: Параметры — отправка значений CFS-параметров + # ═══════════════════════════════════════════════════════ + # send_params_terraform делает 4 вещи: + # 1. Получает cfsParams шаблон операции + # 2. resolveRefSvc — подставляет UUID инстансов для refSvcId-параметров + # 3. Отправляет пользовательские значения + # 4. Отправляет defaults для незаполненных параметров + validate-cfs try: send_params_terraform(client, op_uid, params) except Exception as e: return {"ok": False, "error": str(e), "failed_step": "params", "instance_uid": instance_uid, "op_uid": op_uid, "display_name": display_name} + + # ═══════════════════════════════════════════════════════ + # Шаг 4: POST /instanceOperations/{opUid}/run — запуск + # ═══════════════════════════════════════════════════════ + # После этого операция начинает выполняться в Nubes. + # Дальше — поллинг (poll_until_done), который делает вызывающий код. try: client.post(f"/instanceOperations/{op_uid}/run") except Exception as e: return {"ok": False, "error": str(e), "failed_step": "run", "instance_uid": instance_uid, "op_uid": op_uid, "display_name": display_name} + # Успех — все 4 шага пройдены return {"ok": True, "error": None, "failed_step": None, "instance_uid": instance_uid, "op_uid": op_uid, "display_name": display_name} diff --git a/site/operations/get_instances.py b/site/operations/get_instances.py index c9ad31d..fe22961 100644 --- a/site/operations/get_instances.py +++ b/site/operations/get_instances.py @@ -1,5 +1,10 @@ """ Операции с инстансами Nubes — получение списка и организации через API. + +Используется: + - main.py → отображение инфраструктуры на главной + - api_test.py → список инстансов для выбора + - scenario.py → поиск refSvcId-инстансов """ from api.http_client import HttpClient @@ -7,10 +12,17 @@ from api.http_client import HttpClient def get_instances(client): """GET /instances с пагинацией → ВСЕ инстансы пользователя. - - pageSize=200 — максимум за один запрос. - Остановка по len(batch) < pageSize (НЕ доверяет total — бывали баги). - Возвращает полный list[dict].""" + + pageSize=200 — максимум за один запрос (API лимит). + Остановка по len(batch) < pageSize — НЕ доверяет total из ответа. + + Почему не доверяем total: + Бывали баги где API возвращал total=500, а реально инстансов 3. + Проверка по размеру страницы надёжнее. + + Returns: + list[dict] — каждый dict содержит instanceUid, displayName, serviceId, + svc, explainedStatus, instanceConfigDtCreated и др.""" results = [] page = 1 page_size = 200 @@ -18,8 +30,8 @@ def get_instances(client): data = client.get("/instances", params={"pageSize": page_size, "page": page}) batch = data.get("results", []) or [] results.extend(batch) - # Если страница неполная — данные закончились. - # Не проверяем total, потому что API иногда врёт про total. + # Если страница неполная — данных больше нет. + # Это надёжнее чем проверка page * pageSize >= total. if len(batch) < page_size: break page += 1 @@ -28,8 +40,12 @@ def get_instances(client): def get_organization(client): """Первый инстанс с serviceId=19 (Организация). - - Организация всегда одна на пользователя. Используется в UI: название, статус, ClientID.""" + + Организация всегда одна на пользователя — это его «корневой» инстанс. + Используется в UI: название организации, статус, ClientID. + + Returns: + dict или None — инстанс организации или None если не найден.""" for inst in get_instances(client): if inst.get("serviceId") == 19: return inst diff --git a/site/operations/get_params.py b/site/operations/get_params.py index 4a930df..6555281 100644 --- a/site/operations/get_params.py +++ b/site/operations/get_params.py @@ -1,10 +1,16 @@ """ Операция: получить параметры операции с ТЕКУЩИМИ значениями инстанса. +Зачем: когда пользователь выбирает modify/suspend/resume, мы показываем +ТЕКУЩИЕ значения параметров (из state.params инстанса), а не шаблонные defaults. +Это позволяет видеть что реально настроено, а не то что было при создании. + Источник правды — Nubes API: 1. GET /instances/{instanceUid} → state.params (текущие значения, ключ = код параметра) 2. GET /instanceOperations/default/{opId} (шаблон: код, тип, valueList, dataDescriptor) 3. Смержить: defaultValue = state.params["код"] ?? template.defaultValue + +Также заполняет пустые valueList из state.out (users/databases для PG и др.). """ @@ -12,50 +18,78 @@ def get_params_with_current_values(client, op_id, instance_uid): """ Возвращает список параметров (dict) для заполнения формы операции. - Каждый элемент: - svcOperationCfsParamId — числовой ID параметра + Каждый элемент результирующего списка: + svcOperationCfsParamId — числовой ID параметра (напр. 198) name — код параметра (напр. "durationMs") dataType — тип (напр. "integer >= 0", "boolean", "string") isRequired — обязательный? defaultValue — ТЕКУЩЕЕ значение из state.params или шаблонный default valueList — список допустимых значений (если есть) + refSvcId — ID сервиса для refSvc-параметров dataDescriptor — вложенная структура для map-параметров (если есть) + + Args: + client: HttpClient + op_id: int — svcOperationId (напр. 18 для modify Болванки) + instance_uid: str — UUID инстанса (для получения state.params) + + Returns: + list[dict] — параметры со смерженными текущими значениями """ # --- Шаг 1: текущие значения из инстанса --- + # GET /instances/{uid} возвращает полную информацию об инстансе + # Нас интересует instance.state.params — dict вида: + # {"whereFail": "1", "durationMs": "0", "diskSizeGb": "10"} + # Ключи — СИМВОЛИЧЕСКИЕ имена параметров (не числовые ID!) inst_data = client.get(f"/instances/{instance_uid}") state_params = inst_data.get("instance", {}).get("state", {}).get("params", {}) or {} - # state_params — dict вида {"whereFail": "1", "durationMs": "0", ...} + + # state.out — дополнительные данные (users, databases для PG-сервисов) + # Используется для заполнения пустых valueList (см. ниже) state_out = inst_data.get("instance", {}).get("state", {}).get("out", {}) or {} # --- Шаг 2: шаблон параметров операции --- + # GET /instanceOperations/default/{opId} → svcOperation.cfsParams[] + # Это полный список параметров с типами, default-значениями, valueList и т.д. tmpl_data = client.get(f"/instanceOperations/default/{op_id}") tmpl_params = tmpl_data.get("svcOperation", {}).get("cfsParams", []) or [] - # --- Шаг 3: слияние --- + # --- Шаг 3: слияние шаблона с текущими значениями --- result = [] for p in tmpl_params: pid = p["svcOperationCfsParamId"] # числовой ID (напр. 198) - code = p.get("svcOperationCfsParam", "") # код/имя (напр. "durationMs") + code = p.get("svcOperationCfsParam", "") # символическое имя (напр. "durationMs") data_type = p.get("dataType", "") required = p.get("isRequired", False) tmpl_default = p.get("defaultValue") # шаблонный default value_list = p.get("valueList") - # Заполнить пустые valueList из state.out (пользователи/базы PG и др.) + + # Заполнить ПУСТЫЕ valueList из state.out. + # Некоторые сервисы (PG) возвращают valueList=[] в шаблоне, + # но реальный список значений (пользователи/базы) лежит в state.out. + # Пример: state.out.users = {"pgadmin": {...}, "appuser": {...}} + # → valueList = ["pgadmin", "appuser"] if value_list is not None and not value_list: code_lower = code.lower() + # Параметр типа "пользователь" / "владелец" → state.out.users if ("user" in code_lower or "owner" in code_lower) and isinstance(state_out.get("users"), dict): value_list = sorted(state_out["users"].keys()) + # Параметр типа "база данных" → state.out.databases elif "db" in code_lower and isinstance(state_out.get("databases"), dict): value_list = sorted(state_out["databases"].keys()) - dd = p.get("dataDescriptor") # вложенная структура (map-параметры) - # Берём текущее значение из state.params по коду, - # если нет — используем шаблонный defaultValue. + dd = p.get("dataDescriptor") # вложенная структура (для map-параметров) + + # Выбираем defaultValue: + # Если в state.params есть значение для этого кода → используем его + # Иначе → шаблонный defaultValue current = state_params.get(code) if current is not None: - # Приводим к строке — фронт ожидает строки в полях ввода. - # bool → "true"/"false", dict → JSON, остальное → str() + # Приводим к строке — фронтенд ожидает строки в полях ввода. + # bool → "true"/"false" (не "True"/"False" — JS не поймёт) + # dict → JSON-строка + # остальное → str() if isinstance(current, bool): current = "true" if current else "false" elif isinstance(current, dict): @@ -76,6 +110,7 @@ def get_params_with_current_values(client, op_id, instance_uid): "valueList": _normalize_value_list(value_list), "refSvcId": p.get("refSvcId"), # если параметр ссылается на другой сервис "dataDescriptor": ( + # Для map-параметров: рекурсивно нормализуем вложенные поля { k: { "dataType": v.get("dataType", ""), @@ -93,12 +128,21 @@ def get_params_with_current_values(client, op_id, instance_uid): def _normalize_value_list(vl): - """Nubes API возвращает valueList как CSV-строку или массив. - Приводим к list для фронтенда.""" + """Нормализовать valueList: API может вернуть CSV-строку или массив. + + Приводим к list для фронтенда (всегда массив строк). + + Args: + vl: None | str | list — valueList из API. + + Returns: + list или None — нормализованный список или None если пусто.""" if vl is None: return None + # Уже массив — возвращаем как есть if isinstance(vl, list): return vl + # CSV-строка: "val1,val2,val3" → ["val1", "val2", "val3"] if isinstance(vl, str): return [x.strip() for x in vl.split(',') if x.strip()] or None return None diff --git a/site/operations/get_services.py b/site/operations/get_services.py index c52c563..9457b76 100644 --- a/site/operations/get_services.py +++ b/site/operations/get_services.py @@ -1,5 +1,10 @@ """ Операции с сервисами Nubes — получение списка и деталей через API. + +Используется: + - main.py → список сервисов для UI + - api_test.py → операции сервиса (modify/delete/...) + - scenario.py → поиск svcOperationId для операции """ from api.http_client import HttpClient @@ -7,16 +12,26 @@ from api.http_client import HttpClient def get_services(client): """GET /services → список ВСЕХ сервисов пользователя. - - Возвращает list[dict], каждый с ключами: svcId, svc, svcExtendedName, operations, ...""" + + Возвращает list[dict], каждый с ключами: + svcId, svc (короткое имя), svcExtendedName (полное имя), operations[] + + operations — это ДОСТУПНЫЕ операции, не путать с запущенными.""" data = client.get("/services") return data.get("results", []) def get_service_detail(client, svc_id): """GET /services/{svcId} → детали одного сервиса. - - Возвращает dict с ключами: svc, operations[], ... - operations — список доступных операций (create, modify, delete, ...).""" + + Возвращает dict с ключами: + svc, svcShort, operations[] + + operations — список доступных операций: + [{svcOperationId, operation, isCreate, ...}, ...] + + svcOperationId — ЧИСЛОВОЙ ID операции, нужен для: + - POST /instanceOperations (поле svcOperationId) + - GET /instanceOperations/default/{svcOperationId} (шаблон параметров)""" data = client.get(f"/services/{svc_id}") return data.get("svc", {}) diff --git a/site/operations/poll.py b/site/operations/poll.py index edccfe3..9fbbb8c 100644 --- a/site/operations/poll.py +++ b/site/operations/poll.py @@ -1,9 +1,18 @@ """ -Общий цикл поллинга операции Nubes. +Общий цикл поллинга операции Nubes — ждём завершения или таймаута. -Используется: - - _finish_op (api_test.py) — асинхронно в threading.Thread - - run_scenario (scenario.py) — синхронно в while +Поллинг — это периодический опрос API: + GET /instanceOperations/{opUid}?fields=dtFinish,isSuccessful,errorLog,duration,stages,svc + +Каждые 5 секунд проверяем не появился ли dtFinish (дата завершения). +Как только dtFinish есть → операция завершена → возвращаем результат. + +Используется в двух режимах: + - _finish_op (api_test.py) — асинхронно, в threading.Thread. + Запускается после возврата HTTP-ответа пользователю. + Пользователь видит "RUNNING" и поллит через /api/test/status/. + - run_scenario (scenario.py) — синхронно, в основном потоке сценария. + Сценарий ждёт завершения шага перед переходом к следующему. """ import time @@ -12,23 +21,69 @@ import time def poll_until_done(client, op_uid, timeout=1800): """Поллинг операции до завершения или таймаута. - Возвращает dict: - {status, is_successful, error_log, stages, duration, svc} - status ∈ {OK, FAIL, TIMEOUT} + Алгоритм: + 1. Запоминаем время старта (t0). + 2. В цикле: GET /instanceOperations/{opUid} с нужными полями. + 3. Если dtFinish есть и не пустой → операция завершена. + - isSuccessful=True → status="OK" + - isSuccessful=False → status="FAIL" + 4. Если dtFinish нет → sleep(5) и повторяем. + 5. Если вышли за deadline → status="TIMEOUT". + + Args: + client: HttpClient + op_uid: str — UUID операции (instanceOperationUid) + timeout: int — максимальное время ожидания в секундах. + По умолчанию 1800 (30 минут). + Столько может длиться create сложного сервиса. + + Returns: + dict с ВСЕГДА одинаковыми ключами: + {status, is_successful, error_log, stages, duration, svc} + + status ∈ {"OK", "FAIL", "TIMEOUT"} + OK — операция завершилась успешно + FAIL — операция завершилась с ошибкой + TIMEOUT — не дождались за timeout секунд + + is_successful: bool — True только для OK + error_log: str — текст ошибки из API (или "TIMEOUT after ...") + stages: list — этапы выполнения [{stage, dtStart, dtFinish, isSuccessful, duration}, ...] + duration: float — реальное время поллинга в секундах (округлено до 0.1) + svc: str — имя сервиса (из ответа API) """ + # Фиксируем время начала поллинга t0 = time.time() - deadline = t0 + timeout + deadline = t0 + timeout # крайний срок + + # Основной цикл поллинга while time.time() < deadline: try: + # GET с fields — запрашиваем только нужные поля (экономия трафика) + # dtFinish — ключевое: если есть → операция завершена + # isSuccessful — true/false + # errorLog — текст ошибки если FAIL + # duration — сколько заняла операция по версии API + # stages — этапы (Terraform: plan → apply → ...) + # svc — имя сервиса (может отличаться от того что знаем) data = client.get( f"/instanceOperations/{op_uid}" "?fields=dtFinish,isSuccessful,errorLog,duration,stages,svc" ) except Exception: + # Сетевая ошибка или API временно недоступен — не падаем. + # Ждём 5 секунд и пробуем снова. + # Это нормально: Nubes API иногда отдаёт 502 при высокой нагрузке. time.sleep(5) continue + + # Ответ API: {"instanceOperation": {...}} op = data.get("instanceOperation", {}) dt_finish = op.get("dtFinish") + + # dtFinish может быть None, пустой строкой, или датой. + # Проверяем что dtFinish не пустой: + # str(dt_finish).strip() — отсекаем None→"None" и пустые строки if dt_finish and str(dt_finish).strip(): is_ok = op.get("isSuccessful") return { @@ -39,8 +94,13 @@ def poll_until_done(client, op_uid, timeout=1800): "duration": round(time.time() - t0, 1), "svc": op.get("svc", ""), } + + # Операция ещё не завершена — ждём 5 секунд time.sleep(5) - # Таймаут + + # Вышли из while по deadline — таймаут. + # 30 минут истекли, операция всё ещё выполняется. + # Это редкий случай (обычно create занимает 2-10 минут). return { "status": "TIMEOUT", "is_successful": False, diff --git a/site/operations/scenario.py b/site/operations/scenario.py index f059c47..9900abc 100644 --- a/site/operations/scenario.py +++ b/site/operations/scenario.py @@ -1,14 +1,26 @@ """ -Scenario runner — шаги из БД → API → результат. +Scenario runner — запуск сценария: шаги из БД → API Nubes → результат. -Формат шага (новый): - { "service_id": 1, "operation": "create", "params": {...}, "output": "d1" } - { "service_id": 1, "operation": "modify", "instance_ref": "d1", "params": {...} } - { "service_id": 1, "operation": "delete", "instance_uid": "UUID", "params": {} } +Сценарий — это последовательность операций над сервисами. +Пример: create Болванку → modify параметры → delete Болванку. -Старый формат (без output/instance_ref/instance_uid) — fallback на service_id. +ШАГИ (новый формат с именованными ссылками): + {"service_id": 1, "operation": "create", "params": {...}, "output": "d1"} + — создать инстанс сервиса 1, дать ему имя "d1" для ссылок из следующих шагов + {"service_id": 1, "operation": "modify", "instance_ref": "d1", "params": {...}} + — модифицировать инстанс, созданный в шаге с output="d1" + {"service_id": 1, "operation": "delete", "instance_uid": "UUID", "params": {}} + — удалить конкретный инстанс по UUID -Резолвинг: instance_uid > instance_ref > instance_map[service_id] +Старый формат (без output/instance_ref/instance_uid): + {"service_id": 1, "operation": "create", "params": {...}} + {"service_id": 1, "operation": "modify", "params": {...}} + — fallback на service_id (instance_map[service_id] → последний созданный) + +РЕЗОЛВИНГ instance_uid (приоритет): + 1. Явный instance_uid в шаге + 2. instance_ref → поиск в bindings (output предыдущего create-шага) + 3. instance_map[service_id] → последний инстанс этого сервиса (fallback) """ import json @@ -21,11 +33,28 @@ from operations.poll import poll_until_done from db.pool import get_conn, put_conn from db.save_run import save_run +# Все инстансы созданные сценарием имеют такой префикс в displayName AUTOTEST_PREFIX = "autotest-scenario-" def _resolve_params(cfs_params, symbolic_params): - """Символические коды → {numeric_id: value}.""" + """Конвертировать символические имена параметров → {numeric_id: value}. + + В БД сценарии хранят параметры с СИМВОЛИЧЕСКИМИ именами (напр. "durationMs": "0"). + Но executor.execute_operation требует ЧИСЛОВЫЕ ID (напр. "198": "0"). + + Эта функция делает обратный маппинг: code → svcOperationCfsParamId. + + Args: + cfs_params: list — cfsParams из шаблона операции (содержит svcOperationCfsParam) + symbolic_params: dict — {code: value} из БД + + Returns: + dict — {numeric_id: value} + + Raises: + ValueError если код параметра не найден в шаблоне (опечатка в сценарии).""" + # Строим словарь: code → numeric_id code_to_id = {p.get("svcOperationCfsParam", ""): p["svcOperationCfsParamId"] for p in cfs_params} result = {} for code, val in symbolic_params.items(): @@ -37,7 +66,10 @@ def _resolve_params(cfs_params, symbolic_params): def _save_scenario_run(scenario_run_id, status, current_step=0, duration_sec=None, error_log=None): - """Обновить запись scenario_runs.""" + """Обновить запись scenario_runs в БД. + + Вызывается после КАЖДОГО шага чтобы UI видел прогресс. + При ошибке — пишем в stderr и продолжаем (не роняем сценарий из-за БД).""" conn = get_conn() if not conn: return @@ -58,7 +90,10 @@ def _save_scenario_run(scenario_run_id, status, current_step=0, duration_sec=Non def _update_bindings(scenario_run_id, bindings): - """Сохранить output→instance_uid в scenario_runs.instance_bindings.""" + """Сохранить output→instance_uid в scenario_runs.instance_bindings (JSONB). + + bindings — это словарь: {"d1": "uuid-1", "d2": "uuid-2"} + Нужен для отладки: видеть какие инстансы были созданы на каждом шаге.""" conn = get_conn() if not conn: return @@ -77,40 +112,82 @@ def _update_bindings(scenario_run_id, bindings): def _resolve_instance_uid(step, bindings, instance_map): - """Резолвинг instance_uid по приоритету: явный uid > instance_ref > service_id.""" + """Резолвинг instance_uid по приоритету: явный uid > instance_ref > service_id. + + Args: + step: dict — шаг сценария + bindings: dict — {output_name: instance_uid} (накопленные за сценарий) + instance_map: dict — {service_id: instance_uid} (fallback, последний инстанс сервиса) + + Returns: + str — instance_uid + + Raises: + ValueError если instance_ref не найден в bindings (шаги не в том порядке).""" + # Приоритет 1: явный instance_uid в шаге uid = (step.get("instance_uid") or "").strip() if uid: return uid + + # Приоритет 2: instance_ref → поиск в bindings ref = (step.get("instance_ref") or "").strip() if ref: if ref in bindings: return bindings[ref] raise ValueError(f"instance_ref '{ref}' not found in bindings (step order issue?)") - # Fallback: старый формат по service_id + + # Приоритет 3: fallback — последний инстанс этого сервиса svc_id = step.get("service_id") return instance_map.get(svc_id) def run_scenario(client, steps, client_id, stand, user_email, app_version, scenario_run_id, scenario_name): - """Главный исполнитель сценария. steps — список шагов из БД.""" + """Главный исполнитель сценария — последовательно выполняет шаги. + + Алгоритм для КАЖДОГО шага: + 1. Получить детали сервиса → найти svcOperationId для операции + 2. Получить шаблон параметров → резолвить символические имена в числовые ID + 3. Резолвить instance_uid (create → None, не-create → из bindings/instance_map) + 4. Вызвать execute_operation (единый executor) + 5. Сохранить RUNNING в БД + 6. Поллить через poll_until_done + 7. Сохранить финальный результат (OK/FAIL/TIMEOUT) + instance_meta + 8. Обновить scenario_runs (прогресс для UI) + + Если любой шаг FAIL — сценарий останавливается. + Если все шаги OK — сценарий завершается успешно. + + Args: + client: HttpClient + steps: list[dict] — шаги из БД (scenario_definitions.steps) + client_id: str — ClientID + stand: str — "dev"/"test" + user_email: str — email пользователя + app_version: str — версия приложения + scenario_run_id: int — ID записи в scenario_runs + scenario_name: str — имя сценария + """ total = len(steps) t0 = time.time() - bindings = {} # output_name → instance_uid (новый формат) + + # Две карты для резолвинга instance_uid: + bindings = {} # output_name → instance_uid (новый формат, именованные ссылки) instance_map = {} # service_id → instance_uid (fallback, старый формат) for i, step in enumerate(steps): - step_num = i + 1 + step_num = i + 1 # шаги нумеруются с 1 (для пользователя) svc_id = step.get("service_id") op_name = (step.get("operation") or "").strip() symbolic_params = step.get("params", {}) try: + # Валидация базовых полей шага if not isinstance(svc_id, int) or svc_id <= 0: raise ValueError(f"Invalid service_id: {svc_id!r}") if not op_name: raise ValueError("Empty operation name") - # 1. Операция + # ═══ Шаг 1: получить svcOperationId для операции ═══ detail = get_service_detail(client, svc_id) ops = detail.get("operations", []) op = next((o for o in ops if o.get("operation") == op_name), None) @@ -118,45 +195,51 @@ def run_scenario(client, steps, client_id, stand, user_email, app_version, scena raise ValueError(f"Operation '{op_name}' not found for service {svc_id}") svc_op_id = op["svcOperationId"] - # 2. Параметры + # ═══ Шаг 2: резолвить параметры (символ. имена → числовые ID) ═══ tmpl_data = client.get(f"/instanceOperations/default/{svc_op_id}") cfs_params = tmpl_data.get("svcOperation", {}).get("cfsParams", []) resolved_params = _resolve_params(cfs_params, symbolic_params) - # 3. Резолвинг instance_uid + # ═══ Шаг 3: резолвинг instance_uid ═══ is_create = (op_name == "create") if is_create: + # Для create: instance_uid=None (создаём новый) + # displayName генерируем уникальный: autotest-scenario-{name}-{random6} instance_uid = None display_name = f"{AUTOTEST_PREFIX}{scenario_name}-{uuid.uuid4().hex[:6]}" descr = f"scenario {scenario_name} step {step_num}" else: + # Для не-create: резолвим instance_uid instance_uid = _resolve_instance_uid(step, bindings, instance_map) if not instance_uid: - raise ValueError(f"No instance for step {step_num} — need CREATE before non-create operations") + raise ValueError( + f"No instance for step {step_num} — need CREATE before non-create operations") display_name = None descr = None - # 4. Запуск через executor + # ═══ Шаг 4: запуск через ЕДИНЫЙ executor ═══ + # Это гарантирует идентичное поведение с ручным режимом result = execute_operation( client, svc_id, op_name, instance_uid, resolved_params, svc_op_id=svc_op_id, display_name=display_name, descr=descr, client_id=client_id, stand=stand ) if not result["ok"]: - raise RuntimeError(f"{op_name}: {result['error']} (step: {result['failed_step']})") + raise RuntimeError( + f"{op_name}: {result['error']} (step: {result['failed_step']})") instance_uid = result["instance_uid"] op_uid = result["op_uid"] step_display = result["display_name"] - # Обновить карты - instance_map[svc_id] = instance_uid + # Обновить карты для резолвинга следующих шагов + instance_map[svc_id] = instance_uid # fallback: последний инстанс сервиса output_name = (step.get("output") or "").strip() if is_create and output_name: - bindings[output_name] = instance_uid + bindings[output_name] = instance_uid # именованная ссылка _update_bindings(scenario_run_id, bindings) - # 5. Сохранить RUNNING + # ═══ Шаг 5: сохранить RUNNING в БД (до поллинга) ═══ try: save_run(client_id, stand, user_email, svc_id, "", op_name, svc_op_id, @@ -167,26 +250,27 @@ def run_scenario(client, steps, client_id, stand, user_email, app_version, scena except Exception as e: print(f"[SCENARIO] save_run(RUNNING) failed for step {step_num}: {e}", flush=True) - # 6. Поллинг + # ═══ Шаг 6: поллинг до завершения ═══ poll_result = poll_until_done(client, op_uid) - step_status = poll_result["status"] + step_status = poll_result["status"] # OK / FAIL / TIMEOUT step_error = poll_result["error_log"] step_duration = poll_result["duration"] step_stages = poll_result["stages"] svc_final = poll_result["svc"] - # 7. Сохранить финальный результат + # ═══ Шаг 7: сохранить финальный результат + instance_meta ═══ try: - # Получить полную информацию об инстансе + # Получить полную информацию об инстансе для анализа instance_meta = None try: inst_data = client.get(f"/instances/{instance_uid}") inst = inst_data.get("instance", {}) if isinstance(inst_data, dict) else {} if inst: + # Сохраняем ВСЁ кроме очевидных/избыточных полей instance_meta = {k: v for k, v in inst.items() if k not in ('instanceUid', 'displayName', 'svc', 'serviceId', 'explainedStatus')} except Exception: - pass + pass # не смогли получить meta — не критично save_run(client_id, stand, user_email, svc_id, svc_final, op_name, svc_op_id, op_uid, instance_uid, step_display, @@ -197,17 +281,21 @@ def run_scenario(client, steps, client_id, stand, user_email, app_version, scena except Exception as e: print(f"[SCENARIO] save_run failed for step {step_num}: {e}", flush=True) - # Обновить scenario_runs - _save_scenario_run(scenario_run_id, "RUNNING" if step_num < total else step_status, + # ═══ Шаг 8: обновить scenario_runs (прогресс для UI) ═══ + # Пока не последний шаг — статус RUNNING + _save_scenario_run(scenario_run_id, + "RUNNING" if step_num < total else step_status, step_num, round(time.time() - t0, 1), step_error if step_status != "OK" else None) + # Если шаг FAIL — останавливаем сценарий if step_status != "OK": _save_scenario_run(scenario_run_id, step_status, step_num, round(time.time() - t0, 1), step_error) return except Exception as e: + # Любая ошибка → FAIL, сценарий остановлен import traceback err = f"Step {step_num}: {e}" print(f"[SCENARIO] {err}", flush=True) @@ -216,4 +304,5 @@ def run_scenario(client, steps, client_id, stand, user_email, app_version, scena round(time.time() - t0, 1), err) return + # Все шаги пройдены успешно _save_scenario_run(scenario_run_id, "OK", total, round(time.time() - t0, 1)) diff --git a/site/operations/service_list.py b/site/operations/service_list.py index 93228ab..7999785 100644 --- a/site/operations/service_list.py +++ b/site/operations/service_list.py @@ -1,9 +1,18 @@ """ Чтение списка разрешённых сервисов из services_{stand}.txt. -Формат файла: service_id service_name # описание -Закомментированные (#) строки — исключены. -Пустые строки — игнорируются. +Зачем: не все сервисы должны быть доступны для тестирования. +Файл services_test.txt (или services_dev.txt) содержит только те сервисы, +которые настроены и готовы к автотестам. + +Формат файла: + # Комментарий + 1 dummy # Болванка + 21 vdc # Virtual Data Center + +Закомментированные (#) строки исключаются. +Пустые строки игнорируются. +Первое слово в строке — service_id (число). Файл лежит в site/config/, пушится с кодом, деплоится вместе с приложением. При изменении списка: edit → commit → push → redeploy. @@ -12,23 +21,36 @@ import os import re +# Путь к папке config относительно этого файла: +# service_list.py → site/operations/ +# dirname → site/operations +# dirname(dirname) → site +# + "config" → site/config _CONFIG_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "config") def load_service_ids(stand="test"): - """Вернуть set разрешённых service_id для указанного стенда.""" + """Вернуть set разрешённых service_id для указанного стенда. + + Args: + stand: str — "dev" или "test". + + Returns: + set[int] — множество service_id. Пустой set если файла нет → показываем все.""" path = os.path.join(_CONFIG_DIR, f"services_{stand}.txt") ids = set() try: with open(path) as f: for line in f: line = line.strip() + # Пропускаем пустые строки и комментарии if not line or line.startswith("#"): continue # Формат: "1 dummy # Болванка" + # Берём первое слово, проверяем что это число parts = line.split() if parts and parts[0].isdigit(): ids.add(int(parts[0])) except FileNotFoundError: - pass # файла нет — показываем все сервисы + pass # файла нет — возвращаем пустой set → показываем все сервисы return ids diff --git a/site/operations/terraform.py b/site/operations/terraform.py index 51115de..73a2524 100644 --- a/site/operations/terraform.py +++ b/site/operations/terraform.py @@ -1,21 +1,42 @@ """ -Terraform-совместимые операции: нормализация параметров, refSvc-резолв, отправка. +Terraform-совместимые операции с параметрами: нормализация, refSvc-резолв, отправка. -Используется из api_test.py и scenario.py. +Эти шаги воспроизводят логику Terraform-провайдера Nubes при создании/изменении +инстанса через личный кабинет: + + Шаг 3: Получить cfsParams шаблон операции + Шаг 4: resolveRefSvc — подставить UUID инстансов для refSvcId-параметров + Шаг 5: Отправить пользовательские значения параметров (те что ввёл пользователь) + Шаг 6: Отправить default-значения для незаполненных параметров + Шаг 7: validate-cfs — проверить что все обязательные параметры заполнены + +Используется из executor.py (через send_params_terraform) — и ручной режим, и сценарий. """ import json from operations.get_instances import get_instances +# Параметры, значения которых НЕЛЬЗЯ сохранять в БД в открытом виде. +# При сохранении в runs.params их значения заменяются на "***". REDACT_KEYS = {"password", "secret", "token", "key", "privatekey", "private_key", "passwd", "pass"} def redact_params(params): - """Заменить значения секретных полей на *** перед сохранением в БД.""" + """Заменить значения секретных полей на *** перед сохранением в БД. + + Проверяет КЛЮЧИ параметров (не значения!) на вхождение секретных слов. + Пример: params={"adminPassword": "secret123"} → {"adminPassword": "***"} + + Args: + params: dict — {param_id: value} + + Returns: + dict с замаскированными секретами.""" if not params: return params result = {} for k, v in params.items(): + # Проверяем ключ (имя параметра) на敏感ные слова if any(r in str(k).lower() for r in REDACT_KEYS): result[k] = "***" else: @@ -24,23 +45,45 @@ def redact_params(params): def normalize_value(val, data_type, data_descriptor): - """normalizeUniversalValueV6 — Terraform-equivalent. - Пустые значения нормализуются по типу: integer→"0", boolean→"false", - array→"[]", map/json→"{}", map-fixed→сборка из sub-param defaults.""" + """normalizeUniversalValueV6 — Terraform-эквивалент нормализации значения. + + Что делает: приводит пустые значения к «нулю» соответствующего типа. + Terraform не принимает пустую строку для integer-параметра — нужно "0". + Для map-параметров собирает JSON из sub-param defaults. + + Args: + val: str — значение параметра (может быть "", "null", ...) + data_type: str — тип параметра ("integer >= 0", "boolean", "map(...)", ...) + data_descriptor: dict|None — вложенная структура для map-параметров + + Returns: + str — нормализованное значение.""" + # Очистка: strip + специальные значения null / "" val = (val or "").strip() if val.lower() == "null" or val == '""': val = "" + dt = (data_type or "").lower() + + # Map-параметры с dataDescriptor: + # Если значение пустое — собираем JSON из sub-param defaults. + # Пример: dataDescriptor={"host": {defaultValue: "localhost"}, "port": {defaultValue: "5432"}} + # → "{\"host\":\"localhost\",\"port\":\"5432\"}" if ("map" in dt) and isinstance(data_descriptor, dict) and data_descriptor: if not val or val == "{}": sub_obj = {} for sk, sv in data_descriptor.items(): sd = sv.get("defaultValue") + # default может быть любым типом — приводим к строке для JSON sub_obj[sk] = str(sd) if sd is not None else "" return json.dumps(sub_obj) - return val + return val # значение уже заполнено — не трогаем + + # Значение не пустое — возвращаем как есть if val: return val + + # Пустое значение — подставляем «ноль» типа: if "array" in dt: return "[]" if "map" in dt or "json" in dt: @@ -49,67 +92,126 @@ def normalize_value(val, data_type, data_descriptor): return "0" if "boolean" in dt or "bool" in dt: return "false" + + # string или неизвестный тип — пустая строка return val def resolve_ref_svc(client, cfs_params, params): """resolveRefSvcParamValues — Terraform step 4. - Для параметров с refSvcId: если значение пустое — подставить UUID - существующего инстанса этого сервиса.""" + + Некоторые параметры имеют тип «ссылка на другой сервис» (refSvcId). + Например: параметр «externalIp» ссылается на сервис External IP (svcId=25). + + Если пользователь оставил такой параметр пустым — мы должны найти + существующий инстанс этого сервиса и подставить его UUID. + + Это работает в двух местах: + 1. Верхнеуровневый refSvcId — прямо в параметре + 2. Вложенный refSvcId — внутри dataDescriptor (для map-параметров) + + Args: + client: HttpClient + cfs_params: list — cfsParams из шаблона операции + params: dict — {param_id: value} — МУТИРУЕТСЯ (добавляем UUID) + + Returns: None (мутирует params).""" + # instances загружаем ЛЕНИВО — только если нужен refSvc instances = None + for p in cfs_params: pid = str(p["svcOperationCfsParamId"]) - # Верхнеуровневый refSvcId + + # ── Уровень 1: верхнеуровневый refSvcId ── top_ref = p.get("refSvcId") if top_ref: + # Проверяем: параметр пустой? cur = (params.get(pid) or "").strip() if not cur or cur == '""': + # Загружаем инстансы если ещё не загружены if instances is None: try: instances = get_instances(client) except Exception: - return + return # не можем загрузить — пропускаем + # Ищем первый НЕ deleted инстанс с нужным serviceId match = next((i for i in instances if i.get("serviceId") == top_ref and i.get("explainedStatus") not in ("deleted",)), None) if match: params[pid] = match["instanceUid"] - # Вложенный refSvcId в dataDescriptor + + # ── Уровень 2: вложенный refSvcId в dataDescriptor ── dd = p.get("dataDescriptor") if not isinstance(dd, dict): - continue + continue # нет dataDescriptor — нечего проверять + + # Парсим текущее значение параметра как JSON (map-параметр) user_val = {} if pid in params: try: user_val = json.loads(params[pid]) if params[pid] else {} except (json.JSONDecodeError, ValueError): user_val = {} + + # Проверяем каждый sub-ключ dataDescriptor на refSvcId for sk, sv in dd.items(): ref_id = sv.get("refSvcId") if not ref_id: - continue + continue # не refSvc-параметр + + # Текущее значение sub-ключа cur = user_val.get(sk, "") if cur and cur != '""': - continue + continue # уже заполнено — не трогаем + + # Загружаем инстансы лениво if instances is None: try: instances = get_instances(client) except Exception: return + + # Ищем подходящий инстанс match = next((i for i in instances if i.get("serviceId") == ref_id and i.get("explainedStatus") not in ("deleted",)), None) if match: user_val[sk] = match["instanceUid"] + + # Если были изменения — обновляем значение параметра if user_val: params[pid] = json.dumps(user_val) def send_params_terraform(client, op_uid, params): - """Terraform steps 3-7: cfsParams → resolveRefSvc → send user → send all unsent → validate.""" + """Terraform steps 3-7: cfsParams → resolveRefSvc → send user → send defaults → validate. + + Полный алгоритм отправки параметров операции: + + 3. GET /instanceOperations/{opUid}?fields=cfsParams — получить шаблон + 4. resolveRefSvc — подставить UUID для refSvcId-параметров + 5. Отправить пользовательские значения (те что переданы в params) + 6. Отправить default-значения для НЕзаполненных параметров + (нормализуя через normalize_value) + 7. GET /instanceOperations/{opUid}/validate-cfs — финальная проверка + + Args: + client: HttpClient + op_uid: str — UUID операции + params: dict — {numeric_param_id: string_value} — пользовательские значения + + Raises: + Exception при ошибке валидации (кроме пустого JSON-ответа — это ок).""" + # Шаг 3: получить cfsParams шаблон операции op_det = client.get(f"/instanceOperations/{op_uid}?fields=cfsParams") cfs_params = op_det.get("instanceOperation", {}).get("cfsParams", []) + + # Шаг 4: resolveRefSvc — мутирует params, подставляя UUID resolve_ref_svc(client, cfs_params, params) + + # Шаг 5: отправить пользовательские значения + # sent — множество paramId, чтобы не отправить повторно в шаге 6 sent = set() for pid, pval in params.items(): client.post("/instanceOperationCfsParams", { @@ -118,25 +220,38 @@ def send_params_terraform(client, op_uid, params): "paramValue": str(pval) }) sent.add(int(pid)) + + # Шаг 6: отправить default-значения для незаполненных параметров for p in cfs_params: pid = p["svcOperationCfsParamId"] if pid in sent: - continue + continue # пользователь уже заполнил — не перезаписываем + + # Приоритет значений: paramValue → defaultValue → "" val = p.get("paramValue") if val is None: val = p.get("defaultValue") if val is None: val = "" + + # Нормализуем под тип (пустой integer → "0", map → "{}") val = normalize_value(str(val), p.get("dataType", ""), p.get("dataDescriptor")) + client.post("/instanceOperationCfsParams", { "instanceOperationUid": op_uid, "svcOperationCfsParamId": pid, "paramValue": val }) + + # Шаг 7: validate-cfs — финальная проверка + # Если параметры невалидны — API вернёт ошибку → исключение + # Если ответ пустой или не-JSON — это ОК (значит валидация прошла) 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 + raise # реальная ошибка — пробрасываем выше diff --git a/site/operations/tracker.py b/site/operations/tracker.py index a2ef8b7..4ef90a9 100644 --- a/site/operations/tracker.py +++ b/site/operations/tracker.py @@ -1,20 +1,29 @@ """ -Трекер созданных инстансов. +Трекер созданных инстансов — краткосрочный JSON-файловый кеш. -Хранит UID'ы инстансов, созданных через приложение, в JSON-файлах. -Нужен ТОЛЬКО как краткосрочный fallback — когда инстанс создан, но ещё -не появился в ответе GET /instances (облако может задержать на секунды). +ПРОБЛЕМА: когда мы создаём инстанс через POST /instances, он появляется +в Nubes НЕ МГНОВЕННО. GET /instances может задержать новый инстанс на 5-30 секунд. +В это время пользователь не видит свой инстанс в списке → думает что create не сработал. + +РЕШЕНИЕ: трекер — это JSON-файл в /tmp/, куда мы пишем UID инстанса СРАЗУ после +успешного create. UI сначала ищет инстансы в облаке (cloud-first), а то чего +нет в облаке — добирает из трекера (tracker-fallback). Файлы изолированы по пользователю и стенду: /tmp/instances-{clientId}-{stand}.json -Блокировка fcntl.flock для безопасной работы с несколькими воркерами gunicorn. -При редеплое /tmp/ теряется — это НЕ критично, cloud-first выдача всё покажет. +Блокировка fcntl.flock (эксклюзивная, неблокирующая с retry) — +для безопасной конкурентной работы нескольких gunicorn-воркеров. + +ЖИЗНЕННЫЙ ЦИКЛ: + - add() — executor вызывает сразу после получения instanceUid + - remove() — после успешного delete (в _finish_op и CMDB-delete) + - При редеплое /tmp/ теряется — это НЕ критично (cloud-first всё покажет) Функции: add(client_id, stand, uid, svc_id, name) — добавить инстанс remove(client_id, stand, uid) — удалить (после успешного delete) - list_all(client_id, stand) — все записи списком + list_all(client_id, stand) — все записи списком {svcId, displayName, instanceUid} """ import fcntl import json @@ -23,46 +32,57 @@ import time def _path(client_id, stand): - """Путь к файлу трекера для конкретного пользователя и стенда.""" + """Путь к файлу трекера: /tmp/instances-{clientId}-{stand}.json. + + Изоляция по clientId + stand гарантирует что пользователи не видят чужие инстансы.""" return f"/tmp/instances-{client_id}-{stand}.json" def _acquire_lock(fd): """Эксклюзивная блокировка файла с таймаутом 2 секунды. - + Использует LOCK_NB (неблокирующий) + retry с шагом 50ms. - Если за 2 секунды не взяли лок — возвращает False (не блокируемся навечно).""" + Почему не LOCK_EX с бесконечным ожиданием: + - При подвисании воркера все остальные повиснут навсегда. + - Лучше не записать в трекер, чем уронить запрос. + + Returns: + True — блокировка взята. + False — не смогли за 2 секунды (другой воркер держит слишком долго).""" deadline = time.time() + 2 while True: try: fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) - return True + return True # блокировка взята except BlockingIOError: if time.time() >= deadline: - return False + return False # таймаут — сдаёмся time.sleep(0.05) def _locked_read(path): """Читает JSON-файл под эксклюзивной блокировкой. - - Если файла нет — создаёт и возвращает пустой {}. - Если JSON битый — возвращает {} (переживёт перезапись). - Максимальный размер: 65536 байт.""" + + Если файла нет — создаёт (O_CREAT) и возвращает {}. + Если JSON битый — возвращает {} (переживёт перезапись при следующем add). + Максимальный размер чтения: 65536 байт (64KB) — больше трекеру не нужно. + + Returns: + dict — содержимое файла или {}.""" try: # O_RDWR — чтение+запись, O_CREAT — создать если нет fd = os.open(path, os.O_RDWR | os.O_CREAT, 0o644) except OSError: - return {} # нет прав — молча возвращаем пустой + return {} # нет прав (например /tmp/ только read) — молча try: if not _acquire_lock(fd): - return {} # не смогли заблокировать за 2с — не рискуем + return {} # не смогли заблокировать → не читаем (гонка данных) try: data = os.read(fd, 65536) if data: return json.loads(data.decode("utf-8")) except (json.JSONDecodeError, UnicodeDecodeError): - pass # битый файл — переживём, при следующей записи перезапишется + pass # битый JSON — при следующей записи перезапишется return {} finally: fcntl.flock(fd, fcntl.LOCK_UN) # ВСЕГДА снимаем блокировку @@ -71,28 +91,40 @@ def _locked_read(path): def _locked_write(path, data): """Пишет JSON-файл под эксклюзивной блокировкой. - - Полностью перезаписывает файл: lseek(0) + ftruncate + write.""" + + Полностью перезаписывает файл: lseek(0) + ftruncate + write. + Это атомарно под локом — другие воркеры увидят либо старую, либо новую версию. + + Если не смогли взять лок — НЕ пишем. + Данные не потеряются — инстанс создан в Nubes, cloud-first его подхватит.""" try: fd = os.open(path, os.O_RDWR | os.O_CREAT, 0o644) except OSError: return try: if not _acquire_lock(fd): - return # не смогли заблокировать — не пишем (данные не потеряются, просто не сохранятся) - os.lseek(fd, 0, 0) # в начало файла - os.ftruncate(fd, 0) # обрезать старый контент - os.write(fd, json.dumps(data, indent=2).encode("utf-8")) + return # не смогли заблокировать — пропускаем запись + # Три шага атомарной перезаписи: + os.lseek(fd, 0, 0) # 1. В начало файла + os.ftruncate(fd, 0) # 2. Обрезать старый контент + os.write(fd, json.dumps(data, indent=2).encode("utf-8")) # 3. Записать новый finally: fcntl.flock(fd, fcntl.LOCK_UN) # ВСЕГДА снимаем блокировку os.close(fd) def add(client_id, stand, instance_uid, svc_id, display_name): - """Добавить инстанс в трекер. - - Читает текущий файл → добавляет запись → пишет обратно. - Если инстанс уже есть — перезаписывает (идемпотентно).""" + """Добавить инстанс в трекер (вызывается из executor сразу после create). + + Читает текущий файл → добавляет/обновляет запись → пишет обратно. + Если инстанс уже есть — перезаписывает (идемпотентность). + + Args: + client_id: str — ClientID из JWT + stand: str — "dev"/"test" + instance_uid: str — UUID созданного инстанса + svc_id: int — ID сервиса + display_name: str — displayName инстанса""" p = _path(client_id, stand) data = _locked_read(p) data[instance_uid] = { @@ -105,8 +137,8 @@ def add(client_id, stand, instance_uid, svc_id, display_name): def remove(client_id, stand, instance_uid): """Удалить инстанс из трекера (после успешного delete). - - pop с default=None — не падает если инстанса уже нет.""" + + pop с default=None — не падает если инстанса уже нет (уже удалён ранее).""" p = _path(client_id, stand) data = _locked_read(p) data.pop(instance_uid, None) @@ -114,5 +146,8 @@ def remove(client_id, stand, instance_uid): def list_all(client_id, stand): - """Все записи трекера — список dict'ов {svcId, displayName, instanceUid}.""" + """Все записи трекера — список dict'ов {svcId, displayName, instanceUid}. + + Используется в main.py → api_operations для tracker-fallback. + Конвертирует dict (indexed by UUID) в list.""" return list(_locked_read(_path(client_id, stand)).values())