Загрузка данных
вс кли
self._ws_url = f"wss://{host.rstrip('/')}{WS_Const.WS_HUBS}"
self._negotiate_url = f"https://{host.rstrip('/')}{WS_Const.NEGOTIATE_ENDPOINT}"
self._recv_task: asyncio.Task | None = None
self._ping_task: asyncio.Task | None = None
self._stop_event.set()
if self._ping_task:
self._ping_task.cancel()
try:
await self._ping_task
except asyncio.CancelledError:
pass
raise TimeoutError("Handshake timeout: не получили сообщение с разделителем RS за указанное время")
async def negotiate(self) -> str:
"""
SignalR-negotiate: получает connectionToken для ws-подключения.
connectionToken подставляется в параметр id= URL вебсокета. Без negotiate
api-gateway не регистрирует соединение как полноценного SignalR-клиента,
поэтому он обязателен для корректного подключения.
"""
params = {
"token": self._access_token,
"xUserId": self._x_user_id,
"negotiateVersion": WS_Const.NEGOTIATE_VERSION,
}
response = await asyncio.to_thread(
requests.post,
self._negotiate_url,
params=params,
timeout=WS_Const.NEGOTIATE_TIMEOUT_SECONDS,
)
response.raise_for_status()
connection_token = response.json().get("connectionToken")
if not connection_token:
raise RuntimeError("negotiate не вернул connectionToken")
logger.debug("Получен connectionToken по /negotiate")
return connection_token
attempt += 1
try:
connection_token = await self.negotiate()
self.ws_request = (
f"{self._ws_url}/?token={self._access_token}"
f"&xUserId={self._x_user_id}&id={connection_token}"
)
logger.info(f"Попытка подключения по wss: {self._ws_url}/?token=...&xUserId=...&id=...")
self._ws = await websockets.connect(
await self._handshake()
# Запускаем приём и SignalR keepalive в фоне
self._recv_task = asyncio.create_task(self._recv_loop())
self._ping_task = asyncio.create_task(self._ping_loop())
logger.info("Websocket connected")
continue
async def send_ping(self) -> None:
"""Отправляет SignalR Ping [6] для прикладного keepalive."""
if not self._ws:
return
packet = encode_with_varint_prefix(msgpack.packb([WS_Const.PING_MESSAGE_TYPE], use_bin_type=True))
await self._ws.send(packet)
async def _ping_loop(self) -> None:
"""Периодически шлёт SignalR Ping, пока соединение живо."""
while not self._stop_event.is_set():
await asyncio.sleep(WS_Const.SIGNALR_PING_INTERVAL_SECONDS)
if self._stop_event.is_set() or self._ws is None:
return
try:
await self.send_ping()
except (websockets.ConnectionClosed, OSError):
return
async def invoke(self, target: str, args: list) -> None:
"""
Переподключение websocket соединения
"""
try:
if self._ping_task:
self._ping_task.cancel()
try:
await self._ping_task
except asyncio.CancelledError:
pass
self._ping_task = None
if self._recv_task:
арх
WS_HUBS: str = "/hubs/ldsClientHub"
NEGOTIATE_ENDPOINT: str = "/hubs/ldsClientHub/negotiate"
NEGOTIATE_VERSION: int = 1
NEGOTIATE_TIMEOUT_SECONDS: float = 10.0
CLOSE_TIMEOUT: int = 30
SIGNALR_PING_INTERVAL_SECONDS: float = 2.0 # периодичность SignalR Ping [6]
COMPLETION_MESSAGE_TYPE: int = 3 # Completion
PING_MESSAGE_TYPE: int = 6 # Ping (SignalR keepalive)