Загрузка данных
http_client
import json
import logging
import os
from typing import Optional
import requests
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()
@staticmethod
def make_request(method: str, url: str, **kwargs) -> Optional[requests.Response]:
"""
Обертка для отправки запроса
:param method:
:param url:
:param kwargs: прочие параметры запроса
:return: объект ответа
"""
try:
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:
logging.info(f"[HTTP_CLIENT] Выполняю запрос: METHOD: {method} URL: {url}")
response = self.session.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
@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
@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 _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
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
вс кли
import asyncio
import logging
import time
from datetime import datetime
from typing import Any, Callable, List, Optional
from zoneinfo import ZoneInfo
import msgpack
import websockets
from allure import attach, attachment_type
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 self._should_suppress_recv_attach():
continue
str_message = str(result_message)
logger.info(
f"Обработанное сообщение от api-gateway: {str_message[:200]}... полное сообщение в attach",
)
try:
attach(
str_message,
name=f"Распакованное сообщение от api-gateway {datetime.now(ZoneInfo(WS_Const.ZONE_INFO))}",
attachment_type=attachment_type.TEXT,
)
except (KeyError, RuntimeError) as error:
logger.debug("Allure attach пропущен: %s", error)
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()
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
арх конст
import os
class StandConstants:
MAIN_SUBDOMAIN: str = "web-app"
COMPONENT: str = "lds"
ROOT_DOMAIN: str = "tn.tngrp.ru"
API_GATEWAY_PATH_SEGMENT: str = "lds-api-gateway"
ACKNOWLEDGE_LEAK_URL_PATH: str = "/journals/GetMessages" # TODO Исправить на актуальный в LDS-14845
IMITATE_SIGNAL_URL_PATH: str = "/layerbuilder/ImitateSignal"
GET_BASIC_INFO_URL_PATH: str = "/configurator/GetBasicInfo"
GET_BASIC_INFO_ADMIN_URL_PATH: str = "/configurator/GetBasicInfoAdmin"
GET_TUS_INFORMATION_URL_PATH: str = "/configurator/GetTusInformation"
LAUNCH_LDS_URL_PATH: str = "/configurator/LaunchLds"
STOP_LDS_URL_PATH: str = "/configurator/StopLds"
PING_URL_PATH: str = "/apigateway/Ping"
GET_MESSAGES_URL_PATH: str = "/journals/GetMessages"
GET_OUTPUT_SIGNALS_URL_PATH: str = "/apigateway/GetOutputSignals"
MASK_SIGNAL_URL_PATH: str = "/layerbuilder/MaskSignal"
MASK_LDS_URL_PATH: str = "/core/MaskLds"
UNIMITATE_SIGNAL_URL_PATH: str = "/layerbuilder/UnimitateSignal"
UNMASK_SIGNAL_URL_PATH: str = "/layerbuilder/UnmaskSignal"
UNMASK_LDS_URL_PATH: str = "/core/UnmaskLds"
class ImitatorConstants:
TEST_SETTINGS_KEY_NAME: str = "test_settings"
IMITATOR_FLAGS_KEY_NAME: str = "imitator_flags"
IMITATOR_TIME_FORMAT: str = "%Y%m%dT%H%M%S"
IMITATOR_START_DELAY_S: int = 100
IMITATOR_FINISH_DELAY_MINUTE: float = 2.0
IMITATOR_CHECK_CMD: str = "pgrep -f Playground"
IMITATOR_KILL_CMD: str = "pkill -f Playground"
IMITATOR_PATH = "/data/imitator/lds-flow-playground-csv-latest"
IMITATOR_RUN_CMD: str = f"dotnet {IMITATOR_PATH}/TN.LDS.Flow.Playground.Application.dll"
IMITATOR_LOG_FILE_NAME: str = "imitator.log"
IMITATOR_KEY_NAME: str = "imitator_key"
SERVER_IP_KEY_NAME: str = "server_ip"
SANDBOX_PATH: str = "Sandbox_path" # Удалить при переработке json_config_model.py
SANDBOX_DATA: str = "data"
SANDBOX_RULES: str = "rules.txt"
SANDBOX_TAGS: str = "tags.txt"
STAND_ENV_NAMING: str = os.environ.get("STAND_NAME")[:-1]
CONFIG_PATH: str = f"/data/{STAND_ENV_NAMING}/configs"
SIGNAL_UNIT_CONVERSION_RULES_FILE_NAME: str = "signal_unit_conversion_rules.json"
SIGNAL_UNIT_CONVERSION_RULES_BACKUP_DIR: str = "original_conversion_rules"
SOURCE_TYPE_DEF_VALUE: str = "inflow"
SPEED_DEF_VALUE: int = 1
NS_DEF_VALUE: int = 2
KAFKA_OFFSET_EARLIEST: str = "earliest"
KAFKA_POLL_TIMEOUT_S: float = 1.0
KAFKA_SESSION_TIMEOUT_MS: int = 10000
TEST_ID_KEY: str = "test_id"
AUTOTEST_DATA_PATH: str = "/data/imitator/autotest_data"
POPEN_WAIT_TIMOUT_S: int = 5
LONG_PROCESS_TIMEOUT_S: int = 20
CMD_STATUS_OK: str = "OK"
CMD_STATUS_FAIL: str = "FAIL"
REDIS_STAND_ADDRESS: str = "10.7.49.210"
CORE_START_DELAY_S: int = 5
ENCODING_UTF_8: str = "utf-8"
ENCODING_UTF_8_SIG: str = "utf-8-sig"
ENCODING_LATIN_1: str = "latin-1"
WIN_ENCODING_CP866: str = "cp866" # Нужна только для запуска под WIN
WIN_ENCODING_CP1251: str = "cp1251" # Нужна только для запуска под WIN
OS_NAME_WIN: str = 'nt'
DEFAULT_ENCODINGS = [ENCODING_UTF_8_SIG, ENCODING_UTF_8, WIN_ENCODING_CP866, WIN_ENCODING_CP1251, ENCODING_LATIN_1]
HOST_MAP = {
"dev1": {IMITATOR_KEY_NAME: "DEV1_", SERVER_IP_KEY_NAME: "10.7.49.37"},
"dev2": {IMITATOR_KEY_NAME: "DEV2_", SERVER_IP_KEY_NAME: "10.7.49.38"},
"dev3": {IMITATOR_KEY_NAME: "DEV3_", SERVER_IP_KEY_NAME: "10.7.49.205"},
"test1": {IMITATOR_KEY_NAME: "TEST1_", SERVER_IP_KEY_NAME: "10.7.49.206"},
"test2": {IMITATOR_KEY_NAME: "TEST2_", SERVER_IP_KEY_NAME: "10.7.49.207"},
"test3": {IMITATOR_KEY_NAME: "TEST3_", SERVER_IP_KEY_NAME: "10.7.49.208"},
"test4": {IMITATOR_KEY_NAME: "TEST4_", SERVER_IP_KEY_NAME: "10.7.49.209"},
}
class ClickhouseConstants(ImitatorConstants):
CH_TABLE_NAMES: list = ["lds.records", "lds.records_lastvalue"]
EVO_OBJECT_ID_KEY_NAME: str = "evoObjectId"
EVO_PARAMETER_ID_KEY_NAME: str = "evoParameterId"
OBJECT_ID_KEY_NAME: str = "objectId"
PARAMETER_ID_KEY_NAME: str = "parameterId"
EVO_ID_PAIRS_CHUNK_SIZE: int = 450
NAME_CONTAINER: str = "clickhouse-2"
class DockerConstants:
HOSTNAME_CMD: str = "hostname"
STOP_CMD: str = "docker stop"
START_CMD: str = "docker start"
CHECK_STATUS_CMD: str = "docker inspect -f '{{.State.Status}}'"
RUNNING_STATUS: str = "running"
EXITED_STATUS: str = "exited"
CORE_CONTAINERS_GROUP: list = ["lds-core-node1", "lds-core-node2", "lds-core-node3"]
LB_CONTAINERS_GROUP: list = ["lds-layer-builder-node1", "lds-layer-builder-node2", "lds-layer-builder-node3"]
JOURNAL_CONTAINERS_GROUP: list = ["lds-journals-node1", "lds-journals-node2", "lds-journals-node3"]
WEB_APP_CONTAINERS_GROUP: list = ["lds-web-app-node1", "lds-web-app-node2", "lds-web-app-node3"]
API_GW_CONTAINERS_GROUP: list = ["lds-api-gw-node1", "lds-api-gw-node2", "lds-api-gw-node3"]
REPORTS_CONTAINERS_GROUP: list = ["lds-reports-node1", "lds-reports-node2", "lds-reports-node3"]
class RedisConstants:
LB_REDIS_KEY: str = "lds-layer-builder"
CORE_REDIS_KEY: str = "lds-core"
REDIS_KEY_FIND_CMD: str = "docker exec -i redis-redis-01-1-1 redis-cli KEYS"
REDIS_KEY_DEL_CMD: str = "| xargs -r docker exec -i redis-redis-01-1-1 redis-cli DEL"
class KeycloakClientConstants(StandConstants):
TOKEN_LEEWAY: int = 30
GRANT_TYPE: str = "password"
KEYCLOAK_HEADERS: dict = {"Content-Type": "application/x-www-form-urlencoded"}
TOKEN_KEY: str = "access_token"
TOKEN_URL_PATH: str = "/realms/master/protocol/openid-connect/token"
ISSUED_AT_KEY: str = "issued_at"
EXPIRES_IN_KEY: str = "expires_in"
class TestOpsConstants:
TESTOPS_UPLOAD_ENDPOINT: str = "/upload"
TESTOPS_UPLOAD_ERROR_MSG: str = "Ошибка при загрузке файлов allure отчета"
TESTOPS_UPLOAD_RESPONSE_MSG_KEY: str = "message"
TESTOPS_UPLOAD_FILES_KEY: str = "files"
POST_METHOD: str = "post"
ALLURE_RESULTS_DIR_NAME: str = "allure-results"
GZIP_FILE_SIGNATURE: bytes = b'\x1f\x8b'
class HTTPClientConstants(StandConstants):
GET_METHOD: str = "get"
POST_METHOD: str = "post"
X_SECURITY_SIGNATURE_KEY: str = 'x-security-signature'
X_USER_ID_KEY: str = 'x-user-id'
TESTOPS_UPLOAD_ENDPOINT: str = "/upload"
TESTOPS_ATTACHMENTS_LIST_ENDPOINT: str = "/test_cases/{test_case_id}/attachments"
TESTOPS_LOAD_ATTACHMENT_ENDPOINT: str = "/test_cases/{test_case_id}/attachments/{attachment_id}?download=1"
TESTOPS_ATTACHMENTS_KEY: str = "items"
TESTOPS_ATTACHMENT_FILENAME_KEY: str = "original_filename"
TESTOPS_ATTACHMENT_ID_KEY: str = "id"
TEST_ID_KEY: str = "test_id"
IMITATOR_RUN_DATA_FILENAME: str = "imitator_run_data.tar.gz" # Название архива данных для прогона
DEFAULT_HEADERS: dict = {'content-type': 'application/json'}
STATUS_FORCE_LIST: tuple = (400, 401, 500, 501, 502, 503, 504)
ALLOWED_METHODS: tuple = ("GET", "POST", "PUT", "PATCH", "HEAD")
COLUMNS_SELECTION_DEFAULT: list = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13]
class WebSocketClientConstants(StandConstants):
RS: bytes = b'\x1E' # ASCII Record Separator
HANDSHAKE_WAITING: float | int = 5.0
HANDSHAKE_MESSAGE: str = "{\"protocol\":\"messagepack\",\"version\":1}"
WS_HUBS: str = "/hubs/ldsClientHub"
START_INVOCATION_ID: str = 1
DEFAULT_RECONNECT_INTERVAL: float | int = 5.0
WS_CONNECT_TIMEOUT_SECONDS: float = 120.0
PING_INTERVAL: int = 3
PING_TIMEOUT: int = 5
CLOSE_TIMEOUT: int = 30
DEFAULT_SIGNALR_MESSAGE_TYPE: int = 1 # invocation
STREAM_INVOCATION_MESSAGE_TYPE: int = 4 # StreamInvocation
STREAM_ITEM_MESSAGE_TYPE: int = 2 # StreamItem
COMPLETION_MESSAGE_TYPE: int = 3 # Completion
# Текст ошибки Completion при неуспешном streaming (SignalR CompletionWithDetail)
COMPLETION_ERROR_MESSAGE_INDEX: int = 4
DEFAULT_SIGNALR_MAP_HEADERS: dict = {}
EVENT_TYPE_INDEX = 3
INVOCATION_ID_INDEX = 2
SERVICE_NAME: str = StandConstants.MAIN_SUBDOMAIN
FILTERING_TIMEOUT: int | float = 10.0
ZONE_INFO: str = 'Europe/Moscow'
class MockConstants:
MOCK_DURATION: int = 60
MOCK_TEST_DATA_ID: int = 1
MOCK_TEST_DATA_NAME: str = "mock.tar.gz"
class EnvKeyConstants:
CONNECTION_HOST: str = "CONNECTION_HOST"
KEYCLOAK_URL: str = "KEYCLOAK_URL"
KEYCLOAK_SZI_URL: str = "KEYCLOAK_SZI_URL"
KEYCLOAK_CLIENT_ID: str = "KEYCLOAK_CLIENT_ID"
KEYCLOAK_CLIENT_SECRET: str = "KEYCLOAK_CLIENT_SECRET"
KEYCLOAK_SZI_CLIENT_SECRET: str = "KEYCLOAK_SZI_CLIENT_SECRET"
KEYCLOAK_USERNAME: str = "KEYCLOAK_USERNAME"
KEYCLOAK_PASSWORD: str = "KEYCLOAK_PASSWORD"
TESTOPS_BASE_URL: str = "TESTOPS_BASE_URL"
SSH_KEY_NAME: str = "SSH_KEY_NAME"
SSH_USER_DEV: str = "SSH_USER_DEV"
STAND_NAME: str = "STAND_NAME"
DATA_PATH: str = "DATA_PATH"
OPC_URL: str = "OPC_URL"
TU_ID: str = "TU_ID"