Загрузка данных
ws cli
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, timeout: float = WS_Const.WS_CONNECT_TIMEOUT_SECONDS) -> None:
"""
Цикл подключения с повторными попытками до наступления stop_event
или истечения WS_CONNECT_TIMEOUT_SECONDS.
"""
deadline = time.monotonic() + timeout
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: float = 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
async def reconnect(self):
"""
Переподключение websocket соединения
"""
try:
if self._recv_task:
try:
self._recv_task.cancel()
await self._recv_task
except asyncio.CancelledError:
pass
self._recv_task = None
if self._ws:
await self._ws.close()
self._ws = None
await self._connect_loop(WS_Const.WS_RECONNECT_TIMEOUT_SECONDS)
return self
except (asyncio.TimeoutError, ConnectionError, ConnectionResetError, OSError) as error:
raise RuntimeError(f"Ошибка при попытке переподключения {error}")
arch const
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 = "/core/AcknowledgeLeak"
EXPORT_REPORTS_URL_PATH: str = "/reports/ExportReports"
IMITATE_SIGNAL_URL_PATH: str = "/layerbuilder/ImitateSignal"
GET_BASIC_INFO_URL_PATH: str = "/configurator/GetBasicInfo"
GET_EXPORTED_DATA_LIST_URL_PATH: str = "/journals/GetExportedDataList"
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"
GET_MESSAGES_URL_PATH: str = "/journals/GetMessages"
GET_OUTPUT_SIGNALS_URL_PATH: str = "/apigateway/GetOutputSignals"
LAUNCH_PIG_URL_PATH: str = "/core/LaunchPig"
MASK_SIGNAL_URL_PATH: str = "/layerbuilder/MaskSignal"
MASK_LDS_URL_PATH: str = "/core/MaskLds"
PING_URL_PATH: str = "/apigateway/Ping"
UNIMITATE_SIGNAL_URL_PATH: str = "/layerbuilder/UnimitateSignal"
UNMASK_SIGNAL_URL_PATH: str = "/layerbuilder/UnmaskSignal"
UNMASK_LDS_URL_PATH: str = "/core/UnmaskLds"
ZONE_INFO: str = 'Europe/Moscow'
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"
class DockerConstants:
TIMEOUT_S: int = 20
RETRIES_TIMEOUT_S: int = 5
STOP_CONTAINER_RETRIES: int = 5
HOSTNAME_CMD: str = "hostname"
STOP_CMD: str = "docker stop"
STOP_WITH_TIMEOUT_CMD: str = f"{STOP_CMD} --time {TIMEOUT_S}"
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
WS_RECONNECT_TIMEOUT_SECONDS: float = 20.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 = 30.0
class MockConstants:
MOCK_DURATION: int = 60
MOCK_TEST_DATA_ID: int = 1
MOCK_TEST_DATA_NAME: str = "mock.tar.gz"
MOCK_CHTN_OST_NAME: str = "CHTN"
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"
class VaultConstants:
"""Константы для интеграции с Vault"""
VAULT_ADDR: str = "VAULT_ADDR"
ROLE_ID: str = "LDS_AUTO_VAULT_ROLE"
SECRET_ID: str = "LDS_AUTO_VAULT_SECRET"
VAULT_SKIP_VERIFY: str = "true"
KV_MOUNT: str = "dotnet"
SECRET_PATH_TEMPLATE: str = "config/{environment}/481/{stand_name}/LayerBuilderReportInfoHandler"
ENVIRONMENT_TESTING: str = "Testing"
ENVIRONMENT_DEVELOPMENT: str = "Development"
PROCESS_EMPTY_VALUES_REJECTION_KEY: str = "ProcessEmptyValuesRejection"
VAULT_CMD_TIMEOUT_S: int = 30
enumS
from enum import Enum, IntEnum, IntFlag
from typing import Mapping
class BaseStrEnum(Enum):
def __str__(self) -> str:
return f"{self.name} ({self.value})"
class BaseStrIntFlag(IntFlag):
def __str__(self) -> str:
raw_value = int(self)
if raw_value == 0:
return "0"
active_flags = [f"{flag.name} ({flag.value})" for flag in type(self) if flag.value and flag & self == flag]
if active_flags:
return f"{', '.join(active_flags)}"
return str(raw_value)
class BaseReasonEnum(IntFlag):
"""
.report_text - название на русском языке.
Вывод в формате report_text(value)
"""
def __new__(cls, value: int, report_text: str):
member = int.__new__(cls, value)
member._value_ = value
member.report_text = report_text
return member
def __str__(self) -> str:
raw_value = int(self)
if raw_value == 0:
return "0"
active_flags = [
f"{flag.report_text} ({flag.value})" for flag in type(self) if flag.value and flag & self == flag
]
if active_flags:
return f"{', '.join(active_flags)}"
return str(raw_value)
@classmethod
def report_text_by_value(cls, status_value: int) -> str | None:
"""Текст по числовому значению статуса"""
try:
return cls(status_value).report_text
except ValueError:
return None
class TU(Enum):
"""Legacy-идентификация ТУ: имитатор (tn{id}_tags.txt), tags и конфигурация стенда."""
YAROSLAVL_MOSCOW = (1, "Ярославль - Москва", "volga.json")
TIKHORETSK_NOVOROSSIYSK_2 = (2, "Тихорецк-Новороссийск-2", "tn2.json")
TIKHORETSK_NOVOROSSIYSK_3 = (3, "Тихорецк-Новороссийск-3", "tn3.json")
RODIONOVSKAYA_TIKHORETSKAYA = (4, "Родионовская–Тихорецкая", "lt3_rt.json")
TIKHORETSKAYA_GRUSHEVAYA = (5, "Тихорецкая-6-Грушовая", "tn4_t6g.json")
def __init__(self, tu_id: int, description: str, file_name: str) -> None:
self.id = tu_id
self.description = description
self.file_name = file_name
def __str__(self):
return f"{self.id} - {self.description}"
@classmethod
def get_file_name_by_id(cls, target_id: int) -> str:
for item in cls:
if item.id == target_id:
return item.file_name
raise ValueError(f"ТУ с id = {target_id} не найден")
class AdminTU(Enum):
"""ТУ в Администрировании и WS-подписках"""
TIKHORETSK_NOVOROSSIYSK_3_AUTOTEST = (
"ТН-3-Автотест",
TU.TIKHORETSK_NOVOROSSIYSK_3,
)
TIKHORETSK_NOVOROSSIYSK_3_AUTOTEST_REJECT = (
"ТН-3-Автотест-Отбраковки",
TU.TIKHORETSK_NOVOROSSIYSK_3,
)
TIKHORETSK_NOVOROSSIYSK_3_AUTOTEST_DATA_ABSENCE_FALSE = (
"ТН-3-Автотест-DATA-ABSENCE-FALSE",
TU.TIKHORETSK_NOVOROSSIYSK_3,
)
TIKHORETSK_NOVOROSSIYSK_3_AUTOTEST_DATA_ABSENCE_FALSE_BIK = (
"ТН-3-Автотест-DATA-ABSENCE-BIK",
TU.TIKHORETSK_NOVOROSSIYSK_3,
)
TIKHORETSK_NOVOROSSIYSK_3_AUTOTEST_DATA_ABSENCE_FALSE_SELECT = (
"ТН-3-Автотест-DATA-ABSENCE-SELECT",
TU.TIKHORETSK_NOVOROSSIYSK_3,
)
def __init__(self, admin_name: str, legacy_tu: TU) -> None:
self.admin_name = admin_name
self.legacy_tu = legacy_tu
class ReplyStatus(Enum):
OK = 200
BAD_REQUEST = 400
UNAUTHORIZED = 401
FORBIDDEN = 403
NOT_FOUND = 404
REQUEST_TIMEOUT = 408
CONFLICT = 409
PRECONDITION_FAILED = 412
RANGE_NOT_SATISFIABLE = 416
TOO_MANY_REQUESTS = 429
INTERNAL_SERVER_ERROR = 500
NOT_IMPLEMENTED = 501
SERVICE_UNAVAILABLE = 503
GATEWAY_TIMEOUT = 504
UNKNOWN_ERROR = 520
class ExportedDataType(IntEnum):
"""
Тип экспортируемых данных
"""
STATIONARY_STATUS_REPORT = 6
LDS_STATUS_REPORT = 5
LEAKS_REPORT = 4
REJECTED_REPORT = 7
JOURNAL_REPORT = 1
def to_download_name(self) -> str:
"""Строковый тип для DownloadExportedDataRequest.exportedDataType"""
return _EXPORTED_DATA_TYPE_DOWNLOAD_NAMES[self]
_EXPORTED_DATA_TYPE_DOWNLOAD_NAMES = {
ExportedDataType.STATIONARY_STATUS_REPORT: "MnStateReport",
ExportedDataType.LDS_STATUS_REPORT: "LdsStateReport",
ExportedDataType.LEAKS_REPORT: "LeaksReport",
ExportedDataType.REJECTED_REPORT: "RejectedSignalsReport",
}
class ExportStatus(IntEnum):
"""Статус формирования отчёта в ReportDataExportedNotification.replyContent.exportStatus."""
NOT_READY = 0
DONE = 1
class SouAdminStatus(BaseStrEnum):
"""Статус СОУ в разделе Администрирование (GetBasicInfoAdminResponse)."""
STOPPED = (1, 'СОУ выключена')
RUNNING = (2, 'СОУ включена')
def __new__(cls, value: int, report_text: str) -> "SouAdminStatus":
member = object.__new__(cls)
member._value_ = value
member.report_text = report_text
return member
@classmethod
def report_text_by_value(cls, status_value: int) -> str | None:
"""Текст статуса СОУ в Администрировании по числовому значению."""
try:
return cls(status_value).report_text
except ValueError:
return None
class StationaryStatus(BaseStrEnum):
UNSTATIONARY = (1, 'Нестационарный режим работы МТ')
STATIONARY = (2, 'Стационарный режим работы МТ')
STOPPED = (3, 'МТ в режиме остановленной перекачки')
def __new__(cls, value: int, report_text: str) -> "StationaryStatus":
member = object.__new__(cls)
member._value_ = value
member.report_text = report_text
return member
@classmethod
def report_text_by_value(cls, status_value: int) -> str | None:
"""Текст режима СОУ для отчёта по числовому значению статуса"""
try:
return cls(status_value).report_text
except ValueError:
return None
class LeakStatus(BaseStrEnum):
CONFIRMED = 2
WAITING = 1
POSSIBLE = 3
class LeakLocationStatus(BaseStrEnum):
NODATA = 1 # нет данных
LEFT_FROM_PUMP_STATION = 2 # Слева от МНС
RIGHT_FROM_PUMP_STATION = 3 # Справа от МНС
INSIDE_PUMP_STATION = 4 # Внутри МНС
INSIDE_NPS = 5 # Внутри НПС при неработающей/отсутствующей МНС
class FieldName(Enum):
SECTION_TYPE = "sectionType"
SIGNAL_TYPE = "signalType"
class FilterCriteriaType(Enum):
ONE_OF = "oneOf"
ALL_OF = "allOf"
class FilterCriteriaValue(Enum):
MASK = "mask"
MASK_REASON = "maskReason"
LEAK = "leak"
LEAK_COORDINATE = "leakCoordinate"
PUMPING_STATUS = "pumpingStatus"
FREE_FLOW = "freeFlow"
ACKNOWLEDGE = "acknowledge"
LEAK_TIME = "leakTime"
LEAK_VOLUME = "leakVolume"
LDS_STATUS = "ldsStatus"
CONTROLLED_SITES = "controlledSites"
LINEAR_PARTS = "linearParts"
SERVER_DOWN = "serverDown"
TIME_SYNCHRONIZATION_DISABLE = "timeSynchronizationDisable"
FREE_FLOW_START_COORDINATE = "freeFlowStartCoordinate"
class SortingParam(Enum):
OBJECT_NAME = "objectName"
ADDRESS = "address"
class SortingType(Enum):
ASCENDING = "ascending"
DESCENDING = "descending"
class Direction(Enum):
"""Направление прокрутки"""
PREV = 1
NEXT = 2
FIRST = 3
LAST = 4
class LdsStatus(BaseStrEnum):
"""
Режим работы СОУ.
report_text - значение колонки 'Режим работы СОУ' в xlsx-отчёте об утечках.
"""
FAULTY = (1, "СОУ неисправна")
INITIALIZATION = (2, "СОУ в инициализации")
DEGRADATION = (3, "СОУ в ухудшенных характеристиках")
SERVICEABLE = (4, "СОУ исправна")
def __new__(cls, value: int, report_text: str) -> "LdsStatus":
member = object.__new__(cls)
member._value_ = value
member.report_text = report_text
return member
@classmethod
def report_text_by_value(cls, status_value: int) -> str | None:
"""Текст режима СОУ для отчёта по числовому значению статуса"""
try:
return cls(status_value).report_text
except ValueError:
return None
class ConfirmationStatus(BaseStrEnum):
FAULTY = 0 # Неисправность
AWAITING = 1 # Предварительная
NOT_CONFIRMED = 2 # Не подтверждена
CONFIRMED = 3 # Подтверждена
CONFIRMED_AND_LEAK_CLOSED = 4 # Завершена
class ReservedType(BaseStrEnum):
"""Алгоритмы СОУ"""
FAULTY = 0 # Неисправность
STOP = 1 # Дифференциальный
STATIONARY_FLOW = 2 # Стационарный
UNSTATIONARY_FLOW = 3 # Модельный
BALANCE_IN_NPS = 4 # Баланс внутри НПС
CHANGED_IN_DECISION_MAKING = 5 # Стационарный + Изменено в АПР
CREATED_IN_DECISION_MAKING = 6 # Создано в АПР
class MessageType(BaseStrIntFlag):
USER_ACTION = 0
AUTHENTICATION = 1 # Вход в систему
REJECTION = 1 << 2 # Отбраковка сигналов
LDS_STATUS = 1 << 3 # Режим работы СОУ
INPUT_SIGNALS = 1 << 6 # Входные сигналы
PUMPING_STATUS = 1 << 7 # Режим работы МТ
MASKING_LDS = 1 << 8 # Маскирование СОУ
FREE_FLOWS = 1 << 9 # Самотечное течение
LEAKS = 1 << 10 # Утечка
class MessagePriority(BaseStrIntFlag):
LOW = 1 # Прочее
COMMON = 1 << 1 # Информационное
MEDIUM = 1 << 2 # Значительное
HIGH = 1 << 3 # Важное
VERY_HIGH = 1 << 4 # Особой важности
class LdsStatusDegradation(BaseReasonEnum):
"""
Причины режима работы СОУ: Ухудшение характеристик
"""
LEAK_ON_ADJACENT_DIAGNOSTIC_AREAS = (1 << 0, 'Утечка на соседнем диагностическом участке')
ADDITIVE_INJECTORS_OPERATION = (1 << 1, 'Наличие ПТП')
PIG_SENSOR_PASSAGE = (1 << 2, 'Наличие СОД')
TRIGGERING_EMERGENCY_RESET = (1 << 3, 'Срабатывание аварийного сброса или предохранительных клапанов')
STARTING_PUMPING_OUT_PUMPS = (1 << 4, 'Работа насосов откачки')
EXCEEDING_DISTANCE_BETWEEN_SERVICEABLE_PRESSURE_SENSORS = (
1 << 5,
'Расстояние между ближайшими исправными СИ давления на пути перекачки более 50 км',
)
FAULTY_PRESSURE_SENSORS_AT_PUMP_STATION_NODES = (1 << 6, 'Отказ СИ давления на входе/выходе НПС')
REJECTION_TEMPERATURE_SENSOR = (1 << 7, 'Отказ СИ температуры')
REJECTION_VISCOSITY_SENSOR = (1 << 8, 'Отказ СИ вязкости')
REJECTION_DENSITY_SENSOR = (1 << 9, 'Отказ СИ плотности')
GRAVITY_SECTION_IN_PUMPING_MODE = (1 << 10, 'Наличие самотечного участка/участка с неполным сечением')
ABSENCE_MIN_PRESSURE_SENSORS_REQUIRED_NUMBER = (1 << 11, 'Менее 4 исправных СИ давления на разных КП ЛЧ и НПС')
EXCEEDING_DISTANCE_BETWEEN_FLOW_METERS = (
1 << 12,
'Расстояние между ближайшими исправными СИ расхода на пути перекачки более 200 км',
)
GRAVITY_SECTION_IN_STOPPED_PUMPING_MODE = (
1 << 13,
'Наличие самотечного участка/участка с неполным сечением в режиме остановленной перекачки',
)
class LdsStatusFaulty(BaseReasonEnum):
"""
Причины режима работы СОУ: Неисправность
"""
NO_DATA_SOURCE_CONNECTION = (1 << 0, 'Потеря связи СОУ с СДКУ')
ABSENCE_MIN_PRESSURE_SENSORS_REQUIRED_NUMBER = (1 << 1, 'Менее 4 КП с достоверными СИ давления')
ABSENCE_MIN_FLOW_METERS_REQUIRED_NUMBER = (1 << 2, 'Недостоверность граничного СИ расхода')
class LdsStatusInitialization(BaseReasonEnum):
"""
Причины режима работы СОУ: Инициализация
"""
ACCUMULATION_DATA = (1 << 0, 'Накопление данных')
EXITING_FAULTY_MODE = (1 << 1, 'Выход СОУ из режима «Неисправна»')
COLD_START_OF_SERVERS = (1 << 2, 'Одновременный «холодный» запуск нескольких серверов СОУ')
SWITCHING_SHUT_OFF_IN_STOPPED_PUMPING_MODE = (
1 << 3,
'Переключение запорной арматуры в режиме остановленной перекачки',
)
USER_ACTION = (1 << 4, 'По команде пользователя')
class StationaryReason(BaseReasonEnum):
"""
Причины режима работы МТ: Стационар для ЭФ Журнал
"""
PRESSURE_AND_FLOW_MOVING_AVERAGES_MEET_CRITERIA = (
1 << 0,
'Отклонения давления и расхода не превышают допустимых отклонений',
)
ABSENCE_GRAVITY_SECTION_AND_TECHNOLOGICAL_SWITCHING = (
1 << 1,
'Окончание периода времени после технологических переключений и отсутствия самотечного участка',
)
class UnStationaryReason(BaseReasonEnum):
"""
Причины режима работы МТ: Нестационар
"""
CHANGING_EQUIPMENT_STATUS = (
1 << 0,
'Пуск/остановка трубопровода; включение/отключение магистрального насоса; включение/отключение НПС',
)
CHANGING_WORKING_OF_PUMPING_OUT_PUMPS = (
1 << 1,
'Начало/окончание работы насосов откачки емкостей на НПС и ЛЧ технологического участка',
)
CHANGING_MAIN_PUMPS_ROTATION_SPEED = (
1 << 2,
'Изменение частоты вращения в ручном режиме и/или изменение уставки регулирования '
'в автоматическом режиме работы МНА с ЧРП',
)
CHANGING_BLOCK_VALVES_STATUS = (1 << 3, 'Полное или частичное открытие/закрытие задвижки')
SWITCHING_TANKS = (1 << 4, 'Переключение резервуаров')
CHANGING_ACCEPTANCE_OR_DELIVERY_STATE = (1 << 5, 'Начало или прекращение приема/сдачи нефти/нефтепродуктов')
TRIGGERING_EMERGENCY_RESET_OR_PWSS_OPERATION = (1 << 6, 'Задействование аварийного сброса')
SAFETY_VALVES_ACTUATION = (1 << 7, 'Срабатывание предохранительных клапанов')
CHANGING_PRESSURE_SETTING = (
1 << 8,
'Изменение уставки регулирования по давлению узлов регулирования давления, '
'работающих в автоматическом режиме управления',
)
CHANGING_OPENING_PERCENTAGE_VALVE = (
1 << 9,
'Изменение процента открытия/закрытия заслонки узлов регулирования давления, '
'работающих в ручном режиме управления',
)
CHANGING_ADDITIVE_INJECTOR_STATUS_OR_FLOW = (
1 << 10,
'Начало/окончание ввода ПТП или изменение расхода вводимой ПТП',
)
LEAK_END = (1 << 11, 'Окончание утечки')
TO_OPEN_OR_TO_CLOSE_STATUS = (
1 << 12,
'Наличие сигнала статуса «Открывается»/»Закрывается» запорной арматуры (не в режиме имитации), '
'расположенной в точке, гидравлически связанной с рассматриваемым ДУ',
)
ADJACENT_TU = (
1 << 13,
'Нестационарный режим работы/отсутствие сигнала о режиме работы смежного ТУ, '
'работающего в единой гидравлической системе с защищаемым ТУ',
)
COLD_START = (1 << 14, 'Одновременный «холодный» запуск нескольких серверов СОУ')
class StoppedPumpingReason(BaseReasonEnum):
"""
Причины режима работы МТ: Остановленный
"""
STOPPING_PUMPS = (
1 << 0,
'На ДУ отсутствуют работающие НА, при этом показания СИ расхода не превышают 1 % от максимального значения '
'диапазона измерений всех СИ расхода на технологическом участке',
)
CUTOFF_AREA = (1 << 1, 'Участок отсечен запорной арматурой от подкачек/откачек')
class RejectionCriteria(IntFlag):
"""Критерии отбраковки сигналов criteriaNames"""
QUALITY = 1 << 0 # qualityRejection
RANGE = 1 << 1 # rangeRejection
EMPTY = 1 << 2 # emptyRejection
TIME = 1 << 3 # timeRejection
CONSTANT_SIGNAL = 1 << 4 # constantSignalRejection
DISCHARGE = 1 << 5 # dischargeRejection
VTOR = 1 << 6 # VTORRejection
NEARBY = 1 << 7 # nearbyRejection
DIAGNOSTIC_INFO = 1 << 8 # diagnInfoRejection
@property
def backend_name(self) -> str:
names = {
"QUALITY": "qualityRejection",
"RANGE": "rangeRejection",
"EMPTY": "emptyRejection",
"TIME": "timeRejection",
"CONSTANT_SIGNAL": "constantSignalRejection",
"DISCHARGE": "dischargeRejection",
"SIGMA3": "sigma3Rejection",
"VTOR": "VTORRejection",
"NEARBY": "nearbyRejection",
"DIAGNOSTIC_INFO": "diagnInfoRejection",
}
return names.get(self.name or "", str(int(self)))
def __str__(self) -> str:
raw_value = int(self)
if raw_value == 0:
return "0"
active_flags = [flag.backend_name for flag in type(self) if flag.value and flag & self == flag]
if active_flags:
return f"{'|'.join(active_flags)} ({raw_value})"
return str(raw_value)
class RejectionSensorTag(Enum):
"""
Теги датчиков для тестов отбраковки (id, description=tag)
Айди тегов подставляется на ходу из текущей версии конфы с сервера
"""
KP_8_Pin = (0, "AK.CHTN.LU_TIHVEL.KP_8.SW_8-3.Pin") # nearby_pressure_pin range_upper_pressure range_lower_pres
NPS_TIH_5_Vmom = (0, "AK.CHTN.NPS_TIH_5.UZR_1.Vmom") # diagnostic_info_flowrange_upper_flow range_lower_flow
KP_8_Pout = (0, "AK.CHTN.LU_TIHVEL.KP_8.SW_8-3.Pout") # nearby_pressure_pout
KP_209_1_Pin = (0, "AK.CHTN.LU_VELKRIM.KP_209-1.SW_215-3-1.Pin") # empty_pressure
KP_7_Pin = (0, "AK.CHTN.LU_TIHVEL.KP_7.SW_6-3.Pin") # vtor_pressure
NPS_KRIM_P_Vmom = (0, "AK.CHTN.NPS_KRIM_P.UZR_1.Vmom") # empty_flow quality_flow
def __init__(self, sensor_id: int, description: str) -> None:
self.id = sensor_id
self.description = description
@classmethod
def update_ids_from_config(cls, sensor_ids_by_address: Mapping[str, int]) -> None:
"""
Обновляет sensor_id по tag из конфигурации стенда.
"""
missing_tags = []
for sensor in cls:
sensor_id = sensor_ids_by_address.get(sensor.description)
if sensor_id is None:
missing_tags.append(sensor.description)
continue
sensor.id = sensor_id
if missing_tags:
raise ValueError(f"Не найдены sensor_id для tags: {', '.join(missing_tags)}")
def __str__(self):
return f"{self.id} - {self.description}"
class UserActions(IntFlag):
USER_LOGIN = 1 # Вход пользователя
USER_EXIT = 1 << 1 # Выход пользователя
FAILED_USER_LOGIN = 1 << 2 # Неуспешная попытка входа пользователя
ALGORITHMS_REINITIALIZATION = 1 << 3 # Переинициализация алгоритмов
SIGNAL_MASK_SIM = 1 << 4 # Маскирование и имитация входных сигналов
LDS_MASKING = 1 << 5 # Маскирование СОУ
EXPORT = 1 << 6 # Экспорт и выгрузки
SETTINGS_CHANGE = 1 << 7 # Изменение настроек
LEAK_ACK = 1 << 8 # Квитирование сообщения об утечке
LEAK_REMOVE = 1 << 9 # Исключение неактивных утечек
LDS_ADMIN = 1 << 10 # Администрирование СОУ
PIG_CONTROL = 1 << 11 # Управление СОД
class SiteKpKp(Enum):
"""controlledSiteId, segmentId"""
TIXORECZKAYA_NOVOVELICHKOVSKAYA = (6012, 6013)
NOVOVELICHKOVSKAYA_KRYMSKAYA = (6074, 6075)
KRYMSKAYA_GRUSHOVAYA = (6220, 6221)
BACKUP_ROUTE_BEJSUG = (6242, 6243)
BACKUP_ROUTE_PONURA = (6076, 6077)
BACKUP_ROUTE_KUBAN = (6088, 6089)
NPZ_AFIPSKIJ = (6244, 6245)
NPZ_ILINSKIJ = (6120, 6121)
@property
def controlledSiteId(self) -> int:
return self.value[0]
@property
def segmentId(self) -> int:
return self.value[1]
@property
def controlledSiteSegmentDict(self) -> dict:
return {'controlledSiteId': self.controlledSiteId, 'segmentId': self.segmentId}
@property
def site_key(self) -> tuple[int, int]:
return self.controlledSiteId, self.segmentId
class SignalType(IntFlag):
"""Типы сигналов: режим МТ - 8, режим СОУ - 256, самотеки -16"""
REGLU = 1 << 3
REGSOU = 1 << 8
GRAVITYPIPE = 1 << 4
@property
def backend_name(self) -> str:
names = {
"REGLU": "PumpingStatus",
"REGSOU": "LdsStatus",
"GRAVITYPIPE": "FreeFlow",
}
return names.get(self.name or "", str(int(self)))
def __str__(self) -> str:
raw_value = int(self)
if raw_value == 0:
return "0"
active_flags = [flag.backend_name for flag in type(self) if flag.value and flag & self == flag]
if active_flags:
return f"{'|'.join(active_flags)} ({raw_value})"
return str(raw_value)
class GravityPipe(Enum):
present_gravity = (1, "Наличие самотека")
absent_gravity = (0, "Отсутствие самотека")
def __init__(self, status_id: int, description: str) -> None:
self.id = status_id
self.description = description
def __str__(self):
return f"{self.id} - {self.description}"
class MeasureConversionRule(Enum):
MPA_MEASURE = "MPA_MEASURE"
KG_CM_MEASURE = "KG_CM_MEASURE"
test_const
"""
Общие константы для тестов.
"""
from constants.enums import StationaryStatus
class BaseTN3Constants:
CHTN_OST_NAME = "CHTN"
# ===== Константы для запросов журнала =====
COLUMN_SELECTION_DEF = [
'Time',
'User',
'MainPipeline',
'TechnologicalSection',
'TechnologicalObject',
'ControlPoint',
'Object',
'SignalName',
'Event',
'Value',
'MessageType',
'Tag',
'Status',
]
# ===== Типы сигналов и объектов =====
PRESSURE_SENSOR_OBJECT_TYPE = 2
FLOWMETER_OBJECT_TYPE = 3
PRESSURE_SIGNAL_TYPE = 1
FLOW_SIGNAL_TYPE = 4
# ===== Суффиксы адресов выходных сигналов =====
ADDRESS_SUFFIX_ACK_LEAK = "AckLeak"
ADDRESS_SUFFIX_LEAK = "Leak"
ADDRESS_SUFFIX_MASK = "Mask"
ADDRESS_SUFFIX_POINT_LEAK = "PointLeak"
ADDRESS_SUFFIX_Q_LEAK = "QLeak"
ADDRESS_SUFFIX_TIME_LEAK = "TimeLeak"
ADDRESS_SUFFIX_PUMPING_STATUS = "RegLU"
ADDRESS_SUFFIX_LDS_STATUS = "RegSOU"
# ===== Ключи поиска =====
LEAK_LINEAR_PART_ID_KEY = "id"
CONTROLLED_SITE_ID_AND_SEGMENT_ID = "controlledSiteId"
# ===== Общее количество участков КП-КП =====
LIMIT_CONTROLLED_SITES = 500
COUNT_CONTROLLED_SITES = 114
# ===== Ожидаемые значения выходных сигналов =====
OUTPUT_IS_ACK_LEAK = "1"
OUTPUT_IS_LEAK = "1"
OUTPUT_IS_NOT_LEAK = "0"
OUTPUT_IS_NOT_MASK = "0"
OUTPUT_IS_MASK = "1"
MASS_KG = 3600 # Коэффициент массы, нужно умножить, чтобы получить объем в м3/час
KGS_SM2 = 98066 # Коэффициент давления, нужно умножить, чтобы получить объем в кгс/см2
ALLOWED_VOLUME_DIFF = 0.3 # Относительная погрешность по объему
ALLOWED_DISTANCE_DIFF_METERS = 5000 # Погрешность координаты в метрах
KM_TO_METERS = 1000 # Перевод в метры
LEAK_START_INTERVAL = 2100 # Интервал от старта имитатора до первого обнаружения утечки - 35 минут по умолчанию
LEAK_LOCATION_STATUS = 1
# ===== Параметры выходных сигналов =====
OUTPUT_TEST_DELAY = 120 # Задержка для теста выходных сигналов в секундах
OUTPUT_TIME_FORMAT = "%Y-%m-%dT%H:%M:%SZ" # Формат времени для парсинга выходных сигналов
# ===== Параметры маскирования =====
IS_MASKED_TRUE = True
IS_MASKED_FALSE = False
# ===== Параметры имитации =====
PRESSURE_IMITATION_RANGE = (1, 40)
VOLUME_IMITATION_RANGE = (100, 2400)
GOOD_QUALITY_VAL = 1
# ===== Константы журнала =====
JOURNAL_EVENT_MASK = "Установка признака маскирования"
JOURNAL_EVENT_UNMASK = "Снятие признака маскирования"
JOURNAL_EVENT_IMITATE = "Установка режима имитации сигнала"
JOURNAL_EVENT_UNIMITATE = "Снятие режима имитации сигнала"
JOURNAL_SIGNAL_PRESSURE = "СИ давления"
JOURNAL_SIGNAL_FLOW = "СИ расхода. Расход"
JOURNAL_MESSAGE_TYPE_USER_ACTIONS = "Действия пользователя"
JOURNAL_STATUS_SUCCESS = "Успешно"
JOURNAL_EXPECTED_MSG_COUNT_PER_SIGNAL = 2
JOURNAL_MASK_PAGINATION_LIMIT = 10
JOURNAL_EVENT_POSSIBLE_LEAK = "Возможна утечка"
JOURNAL_EVENT_DETECTED_LEAK = "Утечка."
JOURNAL_MESSAGE_TYPE_LEAKS = "Утечки"
JOURNAL_EVENT_COMPLETED_LEAKS = "Утечка завершена"
JOURNAL_EXPECTED_MASK_MSG_TOTAL = 4
JOURNAL_MASK_EXPECTED_EVENTS = {"Установка признака маскирования", "Снятие признака маскирования"}
JOURNAL_MASK_EXPECTED_SIGNALS = {"Значение давления", "Расход"}
JOURNAL_PAGINATION_LIMIT = 10
JOURNAL_PAGINATION_REJECT_LIMIT = 20
JOURNAL_PAGINATION_STATUS_LIMIT = 120
JOURNAL_STATUS_TOTAL_WAIT = 300 # Время установки режима в данных, в секундах
JOURNAL_EVENT_LEAK_ACKNOWLEDGED = "Сообщение об утечке квитировано"
JOURNAL_EVENT_LDS_INIT_ACCUM_DATA = "СОУ в инициализации (Накопление данных)"
JOURNAL_EVENT_LDS_INIT_COLD_START = "СОУ в инициализации (Одновременный «холодный» запуск нескольких серверов СОУ)"
JOURNAL_MESSAGE_TYPE_LDS_STATUS = "Режим работы СОУ"
JOURNAL_MESSAGE_TYPE_REJECTION = "Отбраковка"
JOURNAL_MESSAGE_EVENT_STATIONARY = (
"Стационарный режим работы МТ (Отклонения давления и расхода не превышают допустимых отклонений)"
)
JOURNAL_MESSAGE_EVENT_NOT_STATIONARY = (
"Нестационарный режим работы МТ (Одновременный «холодный» запуск нескольких серверов СОУ)"
)
JOURNAL_MESSAGE_EVENT_STOP = (
"МТ в режиме остановленной перекачки (На ДУ отсутствуют работающие НА, "
"при этом показания СИ расхода не превышают 1 % отмаксимального значения "
"диапазона измерений всех СИ расхода на технологическом участке)"
)
JOURNAL_TIME_FORMAT = "%Y-%m-%dT%H:%M:%S.%f%z"
SEC_PER_MIN = 60
# ===== Параметры подтверждения =====
IS_ACKNOWLEDGED_FALSE = False
# ===== Параметры BalanceAlgorithmResults =====
BALANCE_ALGORITHM_POLL_INTERVAL = 15 # Интервал опроса подписки в секундах
BALANCE_ALGORITHM_TOTAL_WAIT = 300 # Общее время опроса в секундах
DEBALANCE_TOLERANCE = 0.25 # Допустимое отклонение дебаланса от порога 30%
# ===== Теги датчиков для маскирования и имитации =====
PRESSURE_SENSOR_ADDRESS = "AK.CHTN.LU_TIHVEL.KP_8.SW_8-3.Pout"
FLOWMETER_ADDRESS = "AK.CHTN.NPS_TIH_5.UZR_1.Vmom"
SENSOR_IDS_BY_ADDRESS = {}
# ===== Прочие константы =====
BASIC_MESSAGE_TIMEOUT = 10.0 # Таймаут ожидания сообщений в секундах
SUBSCRIBE_MESSAGE_POLL_ATTEMPTS = 7 # Число чтений из потока подписки до отказа
POLL_BY_TIMEOUT_SECONDS = 5.0 # Таймаут ожидания сообщения в секундах
MASK_MESSAGE_TIMEOUT = 180.0 # Таймаут ожидания сообщений в секундах
PRECISION = 3 # Точность округления для координат
DIGITS_WITH_DOT_PATTERN = r'\d+(?:\.\d+)?' # Регулярное выражение для поиска чисел с точкой
DIAGNOSTIC_AREA_BASE_IDS = {
"Т-Н-3.НПС-5 «Тихорецкая».УЗР СИКН ТН-3 - Т-Н-3.НПС-5 «Тихорецкая».УЗР вых": (9992054907, (9992054908,)),
"Т-Н-3.НПС-5 «Тихорецкая».УЗР вых - Т-Н-3.УЗР НПС-3 «Нововеличковская».": (
9992054908,
(9992054907, 9992054909),
),
"Т-Н-3.УЗР НПС-3 «Нововеличковская». - Т-Н-3.НПС-2 «Крымская».УЗР СИКН Т-К": (
9992054909,
(9992054908, 9992054910, 9992054911),
),
"Т-Н-3.НПС-2 «Крымская».УЗР СИКН Т-К - Т-Н-3.НПС-2 «Крымская».УЗР вых": (9992054910, (9992054909, 9992054912)),
"Н-К.КП-0.УЗР 0км - Н-К.КП-30.УЗР 30км": (9992054911, (9992054909, 9992054913, 9992054914)),
"Т-Н-3.НПС-2 «Крымская».УЗР вых - Т-Н-3.НПС«Крымская».УЗР вых Камеры пуска": (
9992054912,
(9992054910, 9992054915),
),
"Н-К.КП-30.УЗР 30км - ПСП Афипский СИКН 1015 УЗР": (9992054913, (9992054911,)),
"Н-К.УП ИНПЗ.УЗР 4,3км - Н-К.ИНПЗ.УЗР СИКН 1019": (9992054914, (9992054911,)),
"Т-Н-3.НПС«Крымская».УЗР вых Камеры пуска - Т-Н-3.«Грушовая».УЗР-700": (9992054915, (9992054912,)),
}
DIAGNOSTIC_AREA_IDS_EXCLUDED_GRAVITY: list[int] = []
DIAGNOSTIC_AREA_IDS_EXCLUDED_NPS: list[int] = []
REPRESENTATIVE_DIAGNOSTIC_AREA_IDS = [2, 3] # Список показательных ДУ для определения режима СОУ
ZONE_INFO: str = "Europe/Moscow"
SECONDS_PER_HOUR: int = 3600
CRITERIA_NAMES_FIELD: str = 'criteriaNames'
# Словарь вида: {segment.name: (controlledSiteId, segmentId)}, заполняется из файла конфигурации
CONTROLLED_SITE_SEGMENTS = {}
# ===== Ключи для тестовых данных =====
PIPE_ID_KEY: str = "pipe_id"
CONTROL_POINTS_KEY: str = "control_points"
class ExportReportConstants:
"""Константы для теста формирования отчёта об утечках"""
# Максимальное ожидание уведомления о готовности отчёта
NOTIFICATION_TIMEOUT_SECONDS: float = 60.0
# Максимальное время ожидания появления отчёта в списке после уведомления
LIST_POLL_TOTAL_WAIT_SECONDS: float = 60.0
# Интервал между запросами getExportedFilesListRequest
LIST_POLL_INTERVAL_SECONDS: float = 10.0
# Таймаут получения ответа на скачивание
DOWNLOAD_TIMEOUT_SECONDS: float = 60.0
# ===== Имя файла отчёта =====
LEAKS_REPORT_NAME_PART: str = "Отчет об утечках" # подстрока в имени файла/отчёта
XLSX_EXTENSION: str = ".xlsx"
# Сигнатура zip-архива, используется для проверки формата файла по содержимому
ZIP_SIGNATURE: bytes = b'PK\x03\x04'
# ===== Формат даты/времени в отчёте =====
REPORT_DATETIME_FORMAT: str = "%d.%m.%Y %H:%M:%S"
# Регулярное выражение для извлечения двух дат из заголовка
REPORT_HEADER_PERIOD_PATTERN: str = (
r'Отчет об утечках с (?P<period_start>\d{2}\.\d{2}\.\d{4} \d{2}:\d{2}:\d{2})'
r' по (?P<period_end>\d{2}\.\d{2}\.\d{4} \d{2}:\d{2}:\d{2})'
)
# Регулярное выражение для извлечения двух дат из названия файла
REPORT_FILE_NAME_PERIOD_PATTERN: str = (
r'^Отчет об утечках (?P<tu>.+?) '
r'(?P<period_start>\d{2}\.\d{2}\.\d{4} \d{2}_\d{2}_\d{2})'
r' - '
r'(?P<period_end>\d{2}\.\d{2}\.\d{4} \d{2}_\d{2}_\d{2})'
r'\.xlsx$'
)
# Двойная шапка: первая строка - название отчёта с периодом, вторая - названия колонок
REPORT_TITLE_ROW: int = 1
REPORT_COLUMN_HEADERS_ROW: int = 2
REPORT_DATA_FIRST_ROW: int = 3
# ===== Названия колонок =====
COL_DATETIME: str = "Дата и время"
COL_OBJECT: str = "Объект"
COL_LDS_STATUS: str = "Режим работы СОУ"
COL_MASK_INFO: str = "Информация о маскировании"
COL_COORDINATE: str = "Координата"
COL_LEAK_VOLUME: str = "Объемный расход утечки"
COL_MT_MODE: str = "Режим работы МТ"
EXPECTED_COLUMN_HEADERS: list = [
COL_DATETIME,
COL_OBJECT,
COL_LDS_STATUS,
COL_MASK_INFO,
COL_COORDINATE,
COL_LEAK_VOLUME,
COL_MT_MODE,
]
MASKING_NOT_MASKED_TEXT: str = "СОУ не замаскирована"
# ===== Маппинг StationaryStatus <-> текст в колонке "Режим работы МТ" =====
STATIONARY_STATUS_TO_REPORT_TEXT: dict = {
StationaryStatus.UNSTATIONARY.value: "Нестационарный режим работы МТ",
StationaryStatus.STATIONARY.value: "Стационарный режим работы МТ",
StationaryStatus.STOPPED.value: "МТ в режиме остановленной перекачки",
}
# ===== Прочее =====
DEFAULT_SHEET_INDEX: int = 0
SUBSCRIBE_REPORTS_DATA_EXPORTED_REQUEST: str = "SubscribeReportsDataExportedRequest"
EXPORT_REPORTS_COMMAND_REQUEST: str = "ExportReportsCommandRequest"
REPORT_DATA_EXPORTED_NOTIFICATION: str = "ReportDataExportedNotification"
GET_EXPORTED_DATA_LIST_REQUEST: str = "GetExportedDataListRequest"
EXPORTED_DATA_LIST_LIMIT: int = 10
DOWNLOAD_EXPORTED_DATA_REQUEST: str = "DownloadExportedDataRequest"
# Допустимая погрешность при сравнении границ периода отчёта
REPORT_PERIOD_TOLERANCE_MINUTES: int = 1
# Формат даты/времени в имени скачиваемого xlsx-файла
REPORT_FILE_NAME_DATETIME_FORMAT: str = "%d.%m.%Y %H_%M_%S"
class ExportLdsStatusReportConstants:
"""Константы для теста формирования xlsx-отчёта о режиме работы СОУ"""
LDS_STATUS_REPORT_NAME_PART: str = "Отчет о режиме работы СОУ"
SECTION_NAMES: list[str] = [
"НПС-5 Тихорецкая - НПС-3 Нововеличковская",
"НПС-3 Нововеличковская - НПС-2 Крымская",
"НПС-2 Крымская - НПС Грушовая",
]
TOTAL_WORK_DURATION_LABEL: str = "Суммарное время работы:"
ZERO_DURATION_TEXT: str = "0:00:00"
TOTAL_DURATION_TOLERANCE_SECONDS: int = 5
# Число частей времени при split(':') - часы:минуты:секунды (1:02:51) и минуты:секунды (02:51)
DURATION_PARTS_COUNT_H_MM_SS: int = 3
DURATION_PARTS_COUNT_MM_SS: int = 2
REPORT_TITLE_ROW: int = 1
REPORT_COLUMN_HEADERS_ROW: int = 2
REPORT_DATA_FIRST_ROW: int = 3
COL_SECTION: str = "Наименование участка"
COL_FAULTY: str = "Неисправна"
COL_DEGRADATION: str = "В ухудшенных характеристиках"
COL_INITIALIZATION: str = "Инициализация"
COL_SERVICEABLE: str = "Исправна"
MODE_DURATION_COLUMNS: list = [
COL_FAULTY,
COL_DEGRADATION,
COL_INITIALIZATION,
COL_SERVICEABLE,
]
EXPECTED_COLUMN_HEADERS: list = [COL_SECTION, *MODE_DURATION_COLUMNS]
REPORT_HEADER_PERIOD_PATTERN: str = (
r'Отчет о режиме работы СОУ с (?P<period_start>\d{2}\.\d{2}\.\d{4} \d{2}:\d{2}:\d{2})'
r' по (?P<period_end>\d{2}\.\d{2}\.\d{4} \d{2}:\d{2}:\d{2})'
)
REPORT_FILE_NAME_PERIOD_PATTERN: str = (
r'^Отчет о режиме работы СОУ\. (?P<tu>.+?) '
r'(?P<period_start>\d{2}\.\d{2}\.\d{4} \d{2}_\d{2}_\d{2})'
r' - '
r'(?P<period_end>\d{2}\.\d{2}\.\d{4} \d{2}_\d{2}_\d{2})'
r'\.xlsx$'
)
class ExportRejectedReportConstants:
"""Константы для теста формирования xlsx-отчёта об отбракованных входных данных"""
REJECTED_REPORT_NAME_PART: str = "Отчет об отбракованных входных данных"
REJECTED_REPORT_NAME_PART_ALT: str = "Отчёт об отбракованных входных данных"
REPORT_TITLE_ROW: int = 1
REPORT_COLUMN_HEADERS_ROW: int = 2
REPORT_DATA_FIRST_ROW: int = 3
COL_DATETIME: str = "Дата и время"
COL_OBJECT: str = "Объект"
COL_EVENT: str = "Событие"
COL_VALUE: str = "Значение"
COL_DURATION: str = "Продолжительность отбраковки"
COL_TAG: str = "Тег сигнала"
EXPECTED_COLUMN_HEADERS: list = [
COL_DATETIME,
COL_OBJECT,
COL_EVENT,
COL_VALUE,
COL_DURATION,
COL_TAG,
]
REPORT_HEADER_PERIOD_PATTERN: str = (
r'[Оо]тч[её]т об отбракованных входных данных с '
r'(?P<period_start>\d{2}\.\d{2}\.\d{4} \d{2}:\d{2}:\d{2})'
r' по (?P<period_end>\d{2}\.\d{2}\.\d{4} \d{2}:\d{2}:\d{2})'
)
REPORT_FILE_NAME_PERIOD_PATTERN: str = (
r'^[Оо]тч[её]т об отбракованных входных данных (?P<tu>.+?) '
r'(?P<period_start>\d{2}\.\d{2}\.\d{4} \d{2}_\d{2}_\d{2})'
r' - '
r'(?P<period_end>\d{2}\.\d{2}\.\d{4} \d{2}_\d{2}_\d{2})'
r'\.xlsx$'
)
TIME_FILTER_TOLERANCE_SECONDS: int = 60
# Суффикс сигнала в колонке "Объект" отчёта (после последней точки в строке)
REPORT_SIGNAL_FLOW: str = "Расход"
REPORT_SIGNAL_PRESSURE: str = "Давление"
REPORT_SIGNAL_SUFFIX_BY_EXPECTED_NAME: dict = {
BaseTN3Constants.JOURNAL_SIGNAL_FLOW: REPORT_SIGNAL_FLOW,
BaseTN3Constants.JOURNAL_SIGNAL_PRESSURE: REPORT_SIGNAL_PRESSURE,
}
# Разбор колонки "Объект": участок трубопровода и суффикс сигнала разделяются последней точкой
OBJECT_SIGNAL_SEPARATOR: str = "."
OBJECT_SIGNAL_RSPLIT_MAXSPLIT: int = 1
REJECTED_REPORT_HEADER_TITLE_PART: str = "отчет об отбракованных входных данных с"
REJECTED_REPORT_HEADER_TITLE_PART_ALT: str = "отчёт об отбракованных входных данных с"
class MeasureUnitConstants:
MPA_MEASURE: str = "MPa"
KG_CM_MEASURE: str = "kgf/cm^2"
class ExportMtModeReportConstants:
"""Константы для теста формирования xlsx-отчёта о режиме работы МТ"""
MT_MODE_REPORT_NAME_PART: str = "Отчет о режиме работы МТ"
SECTION_NAMES: list[str] = [
"НПС-5 Тихорецкая - НПС-3 Нововеличковская",
"НПС-3 Нововеличковская - НПС-2 Крымская",
"НПС-2 Крымская - НПС Грушовая",
]
TOTAL_WORK_DURATION_LABEL: str = "Суммарное время работы:"
ZERO_DURATION_TEXT: str = "0:00:00"
TOTAL_DURATION_TOLERANCE_SECONDS: int = 5
DURATION_PARTS_COUNT_H_MM_SS: int = 3
DURATION_PARTS_COUNT_MM_SS: int = 2
REPORT_TITLE_ROW: int = 1
REPORT_COLUMN_HEADERS_ROW: int = 2
REPORT_DATA_FIRST_ROW: int = 3
COL_SECTION: str = "Наименование участка"
COL_STOPPED: str = "Остановленный"
COL_UNSTATIONARY: str = "Нестационарный"
COL_STATIONARY: str = "Стационарный"
MODE_DURATION_COLUMNS: list = [
COL_STOPPED,
COL_UNSTATIONARY,
COL_STATIONARY,
]
EXPECTED_COLUMN_HEADERS: list = [COL_SECTION, *MODE_DURATION_COLUMNS]
STATIONARY_STATUS_TO_COLUMN: dict = {
StationaryStatus.STOPPED.value: COL_STOPPED,
StationaryStatus.UNSTATIONARY.value: COL_UNSTATIONARY,
StationaryStatus.STATIONARY.value: COL_STATIONARY,
}
REPORT_HEADER_PERIOD_PATTERN: str = (
r'Отчет о режиме работы МТ с (?P<period_start>\d{2}\.\d{2}\.\d{4} \d{2}:\d{2}:\d{2})'
r' по (?P<period_end>\d{2}\.\d{2}\.\d{4} \d{2}:\d{2}:\d{2})'
)
REPORT_FILE_NAME_PERIOD_PATTERN: str = (
r'^Отчет о режиме работы МТ\. (?P<tu>.+?) '
r'(?P<period_start>\d{2}\.\d{2}\.\d{4} \d{2}_\d{2}_\d{2})'
r' - '
r'(?P<period_end>\d{2}\.\d{2}\.\d{4} \d{2}_\d{2}_\d{2})'
r'\.xlsx$'
)
class LdsConfiguratorConstants:
"""Константы для setup/teardown через раздел Администрирование."""
GET_BASIC_INFO_ADMIN_RETRIES: int = 10
GET_BASIC_INFO_ADMIN_TIMEOUT_SECONDS: float = 2.0
CONFIGURATOR_GET_BASIC_INFO_ADMIN_TIMEOUT_SECONDS: float = 30.0
POLL_TIMEOUT_SECONDS: float = 120.0
VERIFY_UI_SYNC_TIME_SECONDS: float = 300.0
POLL_INTERVAL_SECONDS: float = 15.0
MAIN_PAGE_SYNC_TIMEOUT_SECONDS: float = 30.0
LAUNCHED_AT_TOLERANCE_SECONDS: float = 120.0
GET_BASIC_INFO_ADMIN_REQUEST: str = "GetBasicInfoAdminRequest"
GET_BASIC_INFO_REQUEST: str = "getBasicInfoRequest"
SUBSCRIBE_MAIN_PAGE_INFO_REQUEST: str = "subscribeMainPageInfoRequest"
MAIN_PAGE_INFO_CONTENT: str = "MainPageInfoContent"
STOP_LDS_REQUEST: str = "StopLdsRequest"
LAUNCH_LDS_REQUEST: str = "LaunchLdsRequest"
GET_TUS_INFORMATION_REQUEST: str = "GetTusInformationRequest"
stand set man
import logging
import os
from urllib.parse import urlparse
from clients.subprocess_client import SubprocessClient
from constants.architecture_constants import EnvKeyConstants
from constants.architecture_constants import ImitatorConstants as Im_const
from constants.enums import TU, MeasureConversionRule
from infra.clickhouse_manager import ClickHouseManager
from infra.cmd_generator import ImitatorCmdGenerator
from infra.configuration_manager import ConfigurationManager
from infra.docker_manager import DockerContainerManager
from infra.imitator_data_uploader import ImitatorDataUploader
from infra.imitator_manager import ImitatorManager
from infra.redis_manager import RedisCleaner
from infra.signal_unit_conversion_manager import SignalUnitConversionManager
logger = logging.getLogger(__name__)
class StandSetupManager:
"""
Подготовка стенда к запуску автотестов
Для запуска имитатора:
setup_manager = StandSetupManager(test_duration(minutes), test_data_id, test_data_name, tu_id)
setup_manager.setup_stand_for_imitator_run()
imitator_thread = threading.Thread(target=stand_manager.start_imitator, daemon=True)
core_thread = threading.Thread(target=stand_manager.start_core)
imitator_thread.start()
time.sleep(20)
core_thread.start()
Доступ к времени старта имитатора для расчёта интервалов утечек:
start_time = setup_manager.start_time # datetime объект
"""
def __init__(
self,
duration_m: float, # Максимальное время работы имитатора в минутах
test_data_id: int, # id тест кейса из которого будут загружены данные
test_data_name: str, # Название архива данных имитатора для загрузки из TestOps
tu_id: int,
ost_name: str,
measure_conversion_rules: MeasureConversionRule | None = None,
username: str = os.environ.get(EnvKeyConstants.SSH_USER_DEV),
stand_name: str = os.environ.get(EnvKeyConstants.STAND_NAME),
) -> None:
self._duration_m = duration_m
self._test_data_id = test_data_id
self._test_data_name = test_data_name
self._tu_id = tu_id
self._ost_name = ost_name
self._measure_conversion_rules = measure_conversion_rules
self._username = username
self._stand_name = stand_name
self._configuration_file_name = self._get_configuration_file_name()
self._server_ip = self._get_server_ip() # Получает ip сервера из словаря
self._init_clients()
self._cmd_generator = self._choose_cmd_generator()
self._final_cmd = self._cmd_generator.generate_final_imitator_cmd()
# Экземпляр имитатор менеджера нужно создавать после генерации команды, отдельно от других клиентов
self._imitator_manager = ImitatorManager(self._stand_client, self._final_cmd)
self._remote_data_uploaded = False
@property
def remote_data_uploaded(self) -> bool:
return self._remote_data_uploaded
@property
def start_time(self):
"""
Возвращает время старта имитатора как datetime объект.
Используется для расчёта интервалов утечек в тестах:
- leak_start_time = start_time + LEAK_START_INTERVAL
- leak_end_time = start_time + LEAK_START_INTERVAL + ALLOWED_TIME_DIFF_SECONDS
"""
return self._cmd_generator.start_time
def setup_stand_for_imitator_run(self) -> None:
"""
Обертка, в которой проходит полная подготовка стенда
"""
try:
self._uploader.upload_with_confirm()
self._remote_data_uploaded = True
self.stop_all_containers()
if self._signal_unit_conversion_manager is not None:
self._signal_unit_conversion_manager.setup_signal_unit_conversion_rules()
else:
logger.info(
"[SETUP] [SKIP] measure_conversion_rules не задан в конфигурации набора данных - "
"проверка и настройка единиц измерения (signal_unit_conversion_rules.json) не выполняется"
)
self.clean_redis_and_clickhouse()
self.start_containers_without_core()
except Exception as error:
error_msg = "[SETUP] [ERROR] Ошибка подготовки стенда к запуску имитатора"
logger.exception(error_msg)
raise RuntimeError(error_msg) from error
def restore_signal_unit_conversion_rules(self) -> None:
"""
Возвращает оригинальный signal_unit_conversion_rules.json на стенд.
"""
if self._signal_unit_conversion_manager is None:
return
self._signal_unit_conversion_manager.restore_signal_unit_conversion_rules()
def clean_redis_and_clickhouse(self):
"""
Чистит БД: Clickhouse и Redis
"""
# Копирование файла конфигурации на runner
self._clickhouse_manager.copy_configuration_file_from_stand()
# Чистка ключей Redis
self._redis_cleaner.delete_keys_with_check()
# Чистка ключей ClickHouse
self._clickhouse_manager.delete_clickhouse_keys_with_check()
def get_data_from_configuration(self) -> tuple:
"""
Возвращает словарь address: id из конфигурации, скопированной на runner.
"""
return self._configuration_manager.get_data_from_configuration()
def stop_all_containers(self):
"""
Останавливает все контейнеры и чистит БД: Clickhouse и Redis
"""
self._docker_manager.stop_all_lds_containers()
def start_containers_without_core(self):
"""
Запускает все контейнеры кроме core
"""
# Запуск lds-layer-builder
self._docker_manager.start_lds_layer_builder_containers()
# Запуск lds-journals
self._docker_manager.start_lds_journals_containers()
# Запуск lds-web-app
self._docker_manager.start_lds_web_app_containers()
# Запуск lds-api-gw
self._docker_manager.start_lds_api_gw_containers()
# Запуск lds-reports
self._docker_manager.start_lds_reports_containers()
def start_imitator(self) -> None:
"""
Запускает имитатор и собирает отдельно лог имитатора в файл
"""
try:
self._imitator_manager.run_imitator()
self._imitator_manager.log_imitator_stdout()
except Exception as error:
error_msg = "[SETUP] [ERROR] Ошибка запуска имитатора"
logger.exception(error_msg)
raise RuntimeError(error_msg) from error
def start_core(self) -> None:
"""
Запускает core контейнеры
"""
try:
# Запуск CORE
self._docker_manager.start_lds_core_containers()
except Exception as error:
error_msg = "[SETUP] [ERROR] Ошибка запуска CORE"
logger.exception(error_msg)
raise RuntimeError(error_msg) from error
def stop_imitator_wrapper(self) -> None:
"""
Останавливает имитатор немедленно (без ожидания --stopTime).
В teardown может вызываться даже если имитатор не запущен
"""
try:
if not self._imitator_manager.imitator_process:
logger.info("[TEARDOWN] [SKIP] Имитатор не был запущен")
return
self._imitator_manager.stop_imitator()
logger.info("[TEARDOWN] [OK] Имитатор остановлен")
except Exception as error:
error_msg = "[TEARDOWN] [ERROR] Не удалось остановить имитатор"
logger.exception(error_msg)
raise RuntimeError(error_msg) from error
def server_test_data_remover(self):
"""
Может вызваться в teardown если загрузка не была выполнена
"""
uploader = getattr(self, "_uploader", None)
if uploader is None:
logger.info("[TEARDOWN] [SKIP] uploader не инициализирован, удаление данных со стенда пропущено")
return
if not self._remote_data_uploaded:
logger.info("[TEARDOWN] [SKIP] данные на стенд не загружались, удаление пропущено")
return
try:
uploader.delete_with_confirm()
self._remote_data_uploaded = False
except Exception:
logger.exception("[TEARDOWN] [ALERT] Не удалось удалить данные с сервера")
@staticmethod
def _parse_opc_target() -> tuple[str, int]:
"""
Извлекает хост и порт OPC из переменной окружения OPC_URL.
"""
opc_url = os.environ.get(EnvKeyConstants.OPC_URL)
if not opc_url:
raise RuntimeError(f"[SETUP] [ERROR] Переменная окружения {EnvKeyConstants.OPC_URL} не задана")
parsed = urlparse(opc_url)
if not parsed.hostname or not parsed.port:
raise RuntimeError(
f"[SETUP] [ERROR] Некорректное значение OPC_URL: '{opc_url}'. Ожидается формат вида opc.tcp://host:port"
)
return parsed.hostname, parsed.port
def check_opc_server_status(self, timeout_s: int = 5) -> None:
"""
Проверяет доступность OPC сервера с сервера стенда через /dev/tcp.
"""
host, port = self._parse_opc_target()
check_cmd = (
f"if timeout {timeout_s} bash -lc 'cat < /dev/null > /dev/tcp/{host}/{port}'; "
f"then echo {Im_const.CMD_STATUS_OK}; else echo {Im_const.CMD_STATUS_FAIL}; fi"
)
result = self._stand_client.run_cmd(check_cmd, need_output=True)
if result != Im_const.CMD_STATUS_OK:
raise RuntimeError(f"[SETUP] [ERROR] OPC сервер {host}:{port} недоступен с сервера стенда")
logger.info(f"[SETUP] [OK] OPC сервер {host}:{port} доступен")
def _get_server_ip(self) -> str:
"""
Получает server ip из списка стендов
:return: server ip
"""
try:
return Im_const.HOST_MAP.get(self._stand_name, {}).get(Im_const.SERVER_IP_KEY_NAME)
except Exception as error:
error_msg = f"[SETUP] [ERROR] Не удалось получить server ip для стенда: {self._stand_name}"
logger.exception(error_msg)
raise ValueError(error_msg) from error
def _get_configuration_file_name(self) -> str:
"""
Получает имя файла конфигурации
"""
return TU.get_file_name_by_id(self._tu_id)
def _choose_cmd_generator(self) -> ImitatorCmdGenerator:
"""
Создаёт uploader и генератор команды запуска имитатора по данным из TestOps.
"""
try:
self._uploader = ImitatorDataUploader(
self._stand_client, self._test_data_id, self._test_data_name, self._tu_id
)
self._data_path = self._uploader.remote_temp_dir_path
return ImitatorCmdGenerator(self._data_path, self._ost_name, self._stand_name, self._duration_m)
except Exception as error:
error_msg = "[SETUP] [ERROR] Ошибка при выборе варианта генерации команды запуска имитатора"
logger.exception(error_msg)
raise RuntimeError(error_msg) from error
def _init_clients(self) -> None:
"""
Создает экземпляры необходимых для запуска клиентов
"""
try:
self._stand_client = SubprocessClient(self._username, self._server_ip)
self._infra_client = SubprocessClient(self._username, Im_const.REDIS_STAND_ADDRESS)
self._clickhouse_manager = ClickHouseManager(
self._stand_client, self._infra_client, self._configuration_file_name
)
self._configuration_manager = ConfigurationManager(self._configuration_file_name)
if self._measure_conversion_rules is not None:
self._signal_unit_conversion_manager = SignalUnitConversionManager(
self._stand_client, self._measure_conversion_rules
)
else:
self._signal_unit_conversion_manager = None
self._docker_manager = DockerContainerManager(self._stand_client)
self._redis_cleaner = RedisCleaner(self._infra_client, self._stand_name)
except Exception as error:
error_msg = "[SETUP] [ERROR] Ошибка инициализации клиентов"
logger.exception(error_msg)
raise RuntimeError(error_msg) from error