Загрузка данных


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