Загрузка данных
http
import json
import logging
import os
from datetime import datetime
from typing import Optional
from zoneinfo import ZoneInfo
import requests
from allure import attach, attachment_type
from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry
from constants.architecture_constants import EnvKeyConstants as Env_const
from constants.architecture_constants import HTTPClientConstants as Http_const
logger = logging.getLogger(__name__)
class HttpClient:
"""
Выполняет http запросы
"""
def __init__(self):
self.session = self._create_retry_session()
self.suppress_recv_logging: bool = False
def make_request(self, method: str, url: str, **kwargs) -> Optional[requests.Response]:
"""
Обертка для отправки запроса
:param method:
:param url:
:param kwargs: прочие параметры запроса
:return: объект ответа
"""
try:
if not self.suppress_recv_logging:
logging.info(f"[HTTP_CLIENT] Выполняю запрос: METHOD: {method} URL: {url}")
response = requests.request(method, url, **kwargs)
response.raise_for_status()
return response
except requests.HTTPError as error:
logger.error(
f"[HTTP_CLIENT] [ERROR] При выполнении запроса. METHOD: {method} URL: {url} ERROR_TEXT:{error}"
)
raise
except requests.RequestException:
logger.exception(f"[HTTP_CLIENT] [ERROR] При выполнении запроса. METHOD: {method} URL: {url}")
raise
def make_session_request(self, method: str, url: str, **kwargs) -> Optional[requests.Response]:
"""
Обертка для отправки запроса
:param method:
:param url:
:param kwargs: прочие параметры запроса
:return: объект ответа
"""
try:
if not self.suppress_recv_logging:
logging.info(f"[HTTP_CLIENT] Выполняю запрос: METHOD: {method} URL: {url}")
response = self.session.request(method, url, **kwargs)
response.raise_for_status()
if not self._should_suppress_recv_attach():
self._create_attach_from_http_response(response)
return response
except requests.HTTPError as error:
logger.error(
f"[HTTP_CLIENT] [ERROR] При выполнении запроса. METHOD: {method} URL: {url} ERROR_TEXT:{error}"
)
raise
except requests.RequestException:
logger.exception(f"[HTTP_CLIENT] [ERROR] При выполнении запроса. METHOD: {method} URL: {url}")
raise
@staticmethod
def get_base_url(url_key: str) -> str:
return os.environ.get(url_key)
@staticmethod
def generate_full_url(base_url: str, endpoint: str) -> str:
return f"https://{base_url}{endpoint}"
def post_upload_allure_results(self, files) -> Optional[requests.Response]:
"""
Обертка для выполнения запроса на загрузку отчетов
:return: объект ответа на запрос загрузки отчета
"""
base_testops_url = self.get_base_url(Env_const.TESTOPS_BASE_URL)
full_url = self.generate_full_url(base_testops_url, Http_const.TESTOPS_UPLOAD_ENDPOINT)
response = self.make_request(Http_const.POST_METHOD, full_url, files=files)
return response
def get_attachments_list_by_test_case_id(self, test_case_id: int) -> dict:
"""
Получает список вложений для тест кейса по id через GET запрос к TESTOPS
:return: содержимое ответа на запрос
"""
base_testops_url = self.get_base_url(Env_const.TESTOPS_BASE_URL)
full_endpoint = Http_const.TESTOPS_ATTACHMENTS_LIST_ENDPOINT.format(test_case_id=test_case_id)
full_url = self.generate_full_url(base_testops_url, full_endpoint)
response = self.make_request(Http_const.GET_METHOD, full_url)
return response.json()
def get_test_case_attachment_by_id(self, test_case_id: int, attachment_id: int) -> bytes:
"""
Получает вложение по id через GET запрос к TESTOPS
:return: содержимое ответа на запрос
"""
base_testops_url = self.get_base_url(Env_const.TESTOPS_BASE_URL)
full_endpoint = Http_const.TESTOPS_LOAD_ATTACHMENT_ENDPOINT.format(
test_case_id=test_case_id, attachment_id=attachment_id
)
full_url = self.generate_full_url(base_testops_url, full_endpoint)
response = self.make_request(Http_const.GET_METHOD, full_url)
if not response.content:
logger.exception(f"[HTTP_CLIENT] [ERROR] Пустой response.content при запросе URL: {full_url}")
raise ValueError
return response.content
def _should_suppress_recv_attach(self) -> bool:
if self.suppress_recv_logging:
return True
from utils.helpers import lds_configurator_utils as lds_cfg
return lds_cfg.is_configurator_flow_active()
@staticmethod
def _create_attach_from_http_response(response: requests.Response) -> None:
"""
Создает json вложение http ответа в аллюр
"""
try:
attach(
body=json.dumps(
{"headers": dict(response.headers), "body": response.json()}, ensure_ascii=False, indent=2
),
name=f"HTTP Response от api-gateway {datetime.now(ZoneInfo(Http_const.ZONE_INFO))}",
attachment_type=attachment_type.JSON,
)
except (KeyError, RuntimeError) as error:
logger.debug("Allure attach пропущен: %s", error)
@staticmethod
def _create_retry_session(
retries: int = 3,
backoff_factor: float = 0.3,
status_forcelist: tuple = Http_const.STATUS_FORCE_LIST,
allowed_methods: tuple = Http_const.ALLOWED_METHODS,
) -> requests.Session():
"""
Создает сессию для выполнения нескольких одинаковых запросов
"""
session = requests.Session()
retry = Retry(
total=retries,
read=retries,
connect=retries,
backoff_factor=backoff_factor,
status_forcelist=status_forcelist,
allowed_methods=allowed_methods,
raise_on_status=False,
)
adapter = HTTPAdapter(max_retries=retry)
session.mount("http://", adapter)
session.mount("https://", adapter)
return session
class StandHttpClient(HttpClient):
"""
Выполняет http запросы к стендам
"""
def __init__(self, stand_url: str, token: str) -> None:
super().__init__()
self._stand_url = stand_url
self._token = token
self._headers = self._add_token_to_headers()
self.session.headers.update(self._headers)
def post_request(self, endpoint: str, payload: dict | str) -> Optional[requests.Response]:
"""
Делает http post запрос к стенду
"""
full_url = self.generate_full_url(self._stand_url, endpoint)
json_payload = json.dumps(payload)
response = self.make_session_request(Http_const.POST_METHOD, full_url, data=json_payload)
return response
def _add_token_to_headers(self) -> dict:
"""
Добавляет token в headers запроса
"""
headers = Http_const.DEFAULT_HEADERS.copy()
headers[Http_const.X_SECURITY_SIGNATURE_KEY] = self._token
return headers
keycl
import logging
import time
from typing import Any, Dict, Optional
import requests
from dotenv import load_dotenv
from constants.architecture_constants import KeycloakClientConstants
load_dotenv()
class KeycloakAuthError(Exception):
"""
Исключение для ошибок, возникающих при авторизации в Keycloak.
Атрибуты:
message (str): Сообщение об ошибке.
error_code (int, optional): Код ошибки HTTP.
details (str, optional): Дополнительные детали об ошибке.
"""
def __init__(self, message: Optional[str] = None, error_code: Optional[int] = None, details: Optional[str] = None):
self.message = message or "Ошибка авторизации с Keycloak"
self.error_code = error_code
self.details = details
super().__init__(self.message)
def __str__(self):
"""
Возвращает строковое представление ошибки.
"""
parts = [self.message]
if self.error_code:
parts.append(f"(код: {self.error_code})")
if self.details:
parts.append(f"- детали: {self.details}")
return " ".join(parts)
class KeycloakClient:
"""
Клиент для получения JWT access token из Keycloak с помощью Resource Owner Password Credentials.
"""
def __init__(self, url: str, client_id: str, client_secret: str, username: str, password: str):
"""
Инициализация клиента Keycloak.
Аргументы:
url (str): URL нашего keycloak.
client_id (str): ID клиента автотестов из keycloak.
client_secret (str): Секрет клиента автотестов из keycloak.
username (str): Логин пользователя автотестов из keycloak.
password (str): Пароль пользователя автотестов из keycloak.
"""
self.url = url
self.client_id = client_id
self.client_secret = client_secret
self.username = username
self.password = password
self._token: Optional[str] = None
self._token_data: Optional[Dict[str, Any]]
self.token_leeway = KeycloakClientConstants.TOKEN_LEEWAY
self.grant_type = KeycloakClientConstants.GRANT_TYPE
self.token_key = KeycloakClientConstants.TOKEN_KEY
self.keycloak_headers = KeycloakClientConstants.KEYCLOAK_HEADERS
self.issued_at_key = KeycloakClientConstants.ISSUED_AT_KEY
self.expires_in = KeycloakClientConstants.EXPIRES_IN_KEY
self._validate_creds()
def _validate_creds(self):
required_vars = {
"KEYCLOAK_URL": self.url,
"KEYCLOAK_CLIENT_ID": self.client_id,
"KEYCLOAK_CLIENT_SECRET": self.client_secret,
"KEYCLOAK_USERNAME": self.username,
"KEYCLOAK_PASSWORD": self.password,
}
for name, value in required_vars.items():
if not value:
logging.error(f"Отсутствует обязательная переменная окружения: {name}")
exit(1)
def get_access_token(self):
"""
Получает текущий Access Token. Если токена нет или устарел, запрашивает новый у Keycloak через _request_token()
Возвращает JWT access token.
Рейзим на KeycloakAuthError если не удалось получить токен.
"""
if not self._token or self._is_token_expired():
self._token = self._request_token()
return self._token
def _request_token(self):
"""
Делает запрос к Keycloak для получения Access Token. Возвращает JWT access token.
Если ошибка рейзим на KeycloakAuthError: в случае ошибки авторизации или сети.
"""
data = {
"client_id": self.client_id,
"client_secret": self.client_secret,
"username": self.username,
"password": self.password,
"grant_type": self.grant_type,
}
try:
headers = self.keycloak_headers
response_token = requests.post(self.url, data=data, headers=headers, timeout=5)
response_token.raise_for_status()
token_data = response_token.json()
token = token_data.get(self.token_key)
if not token:
raise KeycloakAuthError("Ответ не содержит access_token", details=str(token_data))
self._token = token
self._token_data = token_data
self._token_data[self.issued_at_key] = int(time.time())
logging.info("Успешно получен access token")
return token_data[self.token_key]
except requests.HTTPError as http_err:
try:
error_json = response_token.json()
error_details = error_json.get("error_description") or error_json.get("error")
except ValueError:
error_details = str(http_err)
error_code = response_token.status_code if response_token is not None else None
logging.error(f"HTTPError при авторизации в Keycloak: {http_err}")
raise KeycloakAuthError(
message="Ошибка HTTP авторизации в Keycloak", error_code=error_code, details=error_details
)
except requests.RequestException as req_err:
logging.error(f"Сетевая ошибка при соединении с Keycloak: {req_err}")
raise KeycloakAuthError(message="Ошибка сетевого соединения с Keycloak", details=str(req_err))
except KeyError as key_err:
logging.error(f"Некорректный формат ответа от Keycloak: {key_err}")
raise KeycloakAuthError(message="Некорректный формат ответа от Keycloak", details=str(key_err))
def _is_token_expired(self):
"""
Проверяет, истёк ли текущий access token.
Для проверки используется поле 'expires_in' в ответе Keycloak
и локальное время получения токена. Если информация отсутствует,
считаем токен недействительным.
"""
if not self._token_data or not self._token:
return True
expires_in = self._token_data.get(self.expires_in)
issued_at = self._token_data.get(self.issued_at_key)
if not expires_in or not issued_at:
return True
now = int(time.time())
if now >= issued_at + expires_in - self.token_leeway:
return True
return False
wscli
import asyncio
import logging
import time
from typing import Any, Callable, List, Optional
import msgpack
import websockets
from websockets import InvalidStatus
from constants.architecture_constants import WebSocketClientConstants as WS_Const
from utils.msgpack_utils.message_filters import is_desired_invocation_id, is_desired_type
from utils.msgpack_utils.msgpack_utils import encode_with_varint_prefix, parse_message
logger = logging.getLogger(__name__)
class WebSocketClient:
"""
Асинхронный ws-клиент для api-gateway по протоколу Async Api
"""
def __init__(
self,
host: str,
access_token: str,
x_user_id: str,
reconnect_interval: float = WS_Const.DEFAULT_RECONNECT_INTERVAL,
):
self._host = host
self._access_token = access_token
self._x_user_id = x_user_id
self._reconnect_interval = reconnect_interval
self._ws_url = f"wss://{host.rstrip('/')}{WS_Const.WS_HUBS}"
self._buffer = b""
self._next_id = WS_Const.START_INVOCATION_ID
self._ws: websockets.ClientConnection | None = None
self.recv_queue: asyncio.Queue[Any] = asyncio.Queue()
self._recv_task: asyncio.Task | None = None
self._stop_event = asyncio.Event()
self._invocation_id: Optional[str] = None
self.suppress_recv_logging: bool = False
@property
def invocation_id(self):
return self._invocation_id
def clear_queue(self):
"""
Очищает очередь путем пересоздания экземпляра класса очереди
"""
while not self.recv_queue.empty():
try:
self.recv_queue.get_nowait()
self.recv_queue.task_done()
except asyncio.QueueEmpty as message_empty:
logger.info(f"Очередь сообщений очищена: {message_empty}")
break
async def __aenter__(self):
await self._connect_loop()
return self
async def __aexit__(self, exc_type, exc, tb):
self._stop_event.set()
if self._ws:
await self._ws.close()
if self._recv_task:
await self._recv_task
# TODO: дергать ручку завершения сессии в LDS-4083
async def _handshake(self) -> None:
payload = WS_Const.HANDSHAKE_MESSAGE + WS_Const.RS.decode()
logger.debug(
f"Отправлен handshake: {payload}",
)
await self._ws.send(payload)
buf = self._buffer
finish = time.monotonic() + WS_Const.HANDSHAKE_WAITING
while time.monotonic() < finish:
chunk = await self._ws.recv()
logger.debug(f"Ответ на handshake: {chunk}")
buf += chunk.encode() if isinstance(chunk, str) else chunk
if WS_Const.RS in buf:
return
raise TimeoutError("Handshake timeout: не получили сообщение с разделителем RS за указанное время")
async def _connect_loop(self) -> None:
"""
Цикл подключения с повторными попытками до наступления stop_event
или истечения WS_CONNECT_TIMEOUT_SECONDS.
"""
deadline = time.monotonic() + WS_Const.WS_CONNECT_TIMEOUT_SECONDS
attempt = 0
transient_errors = (ConnectionError, OSError, asyncio.TimeoutError, InvalidStatus)
while not self._stop_event.is_set():
attempt += 1
try:
self.ws_request = f"{self._ws_url}/?token={self._access_token}&xUserId={self._x_user_id}"
logger.info(f"Попытка подключения по wss: {self._ws_url}/?token=...&xUserId=...")
self._ws = await websockets.connect(
self.ws_request,
ping_interval=WS_Const.PING_INTERVAL,
ping_timeout=WS_Const.PING_TIMEOUT,
close_timeout=WS_Const.CLOSE_TIMEOUT,
)
# Handshake
await self._handshake()
# Запускаем приём в фоне
self._recv_task = asyncio.create_task(self._recv_loop())
logger.info("Websocket connected")
return
except transient_errors as exc:
remaining = deadline - time.monotonic()
if remaining <= 0:
raise TimeoutError(
f"WSS hub не готов за {WS_Const.WS_CONNECT_TIMEOUT_SECONDS} с " f"(попыток: {attempt}): {exc}"
) from exc
status_info = ""
if isinstance(exc, InvalidStatus):
response = getattr(exc, "response", None)
if response is not None:
status_info = f", HTTP {response.status_code}"
logger.warning(
"WSS подключение не установлено (попытка %s%s): %s. " "Повтор через %s с, осталось %.0f с",
attempt,
status_info,
exc,
self._reconnect_interval,
remaining,
)
await asyncio.sleep(self._reconnect_interval)
async def _recv_loop(self) -> None:
"""
Прием сообщений, парсинг и отправка в очередь.
"""
assert self._ws is not None
while not self._stop_event.is_set():
try:
chunk = await self._ws.recv()
if not self.suppress_recv_logging:
logger.debug(f"Сырые биты до обработки: {chunk[:100]}")
except websockets.ConnectionClosed as e:
logger.warning(f"WebSocket соединение разорвано: {e}")
return
result_message = parse_message(chunk)
await self.recv_queue.put(result_message)
if not self.suppress_recv_logging:
str_message = str(result_message)
logger.info(
f"Обработанное сообщение от api-gateway: {str_message[:200]}... полное сообщение в attach",
)
continue
async def invoke(self, target: str, args: list) -> None:
"""
Отправляет удаленный вызов invocation о websocket соединению
Метод формирует сообщение по протоколу SignalR, включая:
- типа сообщения
- заголовки
- уникальный идентификатор запроса
- имя целевого метода
- аргументы запроса
Сообщение запаковывается в messagepack и отправляется через текущее websocket соединение
"""
if not self._ws:
raise websockets.WebSocketException("Не установлено подключение по wss")
self._invocation_id = str(self._next_id)
self._next_id += 1
invocation = [
WS_Const.DEFAULT_SIGNALR_MESSAGE_TYPE,
WS_Const.DEFAULT_SIGNALR_MAP_HEADERS,
self._invocation_id,
target,
[args],
]
logger.info(f"Сообщение подготовлено к отправке: {invocation}")
payload = msgpack.packb(invocation, use_bin_type=True)
packet = encode_with_varint_prefix(payload)
logger.debug(f"Отправляем сообщение: {packet}")
await self._ws.send(packet)
async def invoke_stream(self, target: str, args: list) -> None:
"""
Отправляет streaming-вызов (StreamInvocation) по протоколу SignalR.
"""
if not self._ws:
raise websockets.WebSocketException("Не установлено подключение по wss")
self._invocation_id = str(self._next_id)
self._next_id += 1
invocation = [
WS_Const.STREAM_INVOCATION_MESSAGE_TYPE,
WS_Const.DEFAULT_SIGNALR_MAP_HEADERS,
self._invocation_id,
target,
[args],
]
logger.info(f"Streaming-сообщение подготовлено к отправке: {invocation}")
payload = msgpack.packb(invocation, use_bin_type=True)
packet = encode_with_varint_prefix(payload)
await self._ws.send(packet)
async def receive_by_type(self, message_type: str, timeout: WS_Const.FILTERING_TIMEOUT) -> List[Any]:
"""
Фильтрует сообщения по message_type
"""
try:
return await self._receive_by(filter_func=lambda msg: is_desired_type(msg, message_type), timeout=timeout)
except websockets.WebSocketException:
raise websockets.WebSocketException(f"Ошибка при фильтрации сообщений по {message_type}")
async def receive_by_invocation_id(
self, invocation_id: str, timeout: float = WS_Const.FILTERING_TIMEOUT
) -> List[Any]:
"""
Фильтрует сообщения по invocation_id
"""
try:
return await self._receive_by(
filter_func=lambda msg: is_desired_invocation_id(msg, invocation_id), timeout=timeout
)
except websockets.WebSocketException:
raise websockets.WebSocketException("Ошибка при фильтрации сообщений по invocation_id")
async def _receive_by(self, filter_func: Callable[[list], bool], timeout: float) -> List[Any]:
"""
Ждет и фильтрует сообщение по filter_func
"""
# 1) Единая точка вычисления дедлайна
deadline = time.monotonic() + timeout
while True:
# 2) Остаток времени до таймаута
remaining = deadline - time.monotonic()
if remaining <= 0:
raise asyncio.TimeoutError(f"Timeout при фильтрации сообщений {timeout:.1f} секунд")
try:
# 3) Получает сообщение
msg = await asyncio.wait_for(self.recv_queue.get(), timeout=remaining)
except asyncio.TimeoutError:
# 4) Явно перехватывает и пробрасывает свой TimeoutError
raise asyncio.TimeoutError(f"Timeout при фильтрации сообщений {timeout:.1f} секунд")
# 5) Фильтрация по filter_func
if isinstance(msg, list) and filter_func(msg):
return msg