diff --git a/app/scripts/check_trades_websocket.py b/app/scripts/check_trades_websocket.py new file mode 100644 index 0000000..25bdb30 --- /dev/null +++ b/app/scripts/check_trades_websocket.py @@ -0,0 +1,269 @@ +# app/scripts/check_trades_websocket.py + +from __future__ import annotations + +import asyncio +import json +from typing import Any +from uuid import uuid4 + +import websockets +from websockets.typing import Subprotocol + +from src.core.config import load_settings + + +SYMBOLS = ( + "BTC/USD_LEVERAGE", +) + +MAX_MESSAGES = 100 +TOTAL_TIMEOUT_SECONDS = 3000.0 +RECEIVE_TIMEOUT_SECONDS = 15.0 + + +def build_ws_url() -> str: + settings = load_settings() + + raw_url = settings.exchange_ws_url or settings.exchange_base_url + + if raw_url.startswith("https://"): + raw_url = raw_url.replace("https://", "wss://", 1) + elif raw_url.startswith("http://"): + raw_url = raw_url.replace("http://", "ws://", 1) + + raw_url = raw_url.rstrip("/") + + if raw_url.endswith("/connect"): + return raw_url + + return f"{raw_url}/connect" + + +def build_headers() -> dict[str, str]: + settings = load_settings() + + headers = { + "Origin": settings.exchange_base_url.rstrip("/"), + "Content-Type": "application/json", + } + + if settings.exchange_api_key: + headers["X-MBX-APIKEY"] = settings.exchange_api_key + + return headers + + +def build_subscribe_request( + symbols: tuple[str, ...], +) -> dict[str, Any]: + return { + "correlationId": str(uuid4()), + "destination": "trades.subscribe", + "payload": { + "symbols": list(symbols), + }, + } + + +def format_json( + raw_message: str | bytes | bytearray | memoryview, +) -> str: + if isinstance(raw_message, str): + text = raw_message + elif isinstance(raw_message, memoryview): + text = raw_message.tobytes().decode( + "utf-8", + errors="replace", + ) + else: + text = bytes(raw_message).decode( + "utf-8", + errors="replace", + ) + + try: + document = json.loads(text) + except json.JSONDecodeError: + return text + + return json.dumps( + document, + ensure_ascii=False, + indent=2, + sort_keys=True, + ) + + +def decode_document( + raw_message: str | bytes | bytearray | memoryview, +) -> object | None: + if isinstance(raw_message, str): + text = raw_message + elif isinstance(raw_message, memoryview): + text = raw_message.tobytes().decode( + "utf-8", + errors="replace", + ) + else: + text = bytes(raw_message).decode( + "utf-8", + errors="replace", + ) + + try: + return json.loads(text) + except json.JSONDecodeError: + return None + + +def is_trade_event(document: object) -> bool: + if not isinstance(document, dict): + return False + + return document.get("destination") == "internal.trade" + + +def elapsed_seconds( + loop: asyncio.AbstractEventLoop, + started_at: float, +) -> float: + return max(0.0, loop.time() - started_at) + + +async def receive_trades() -> None: + settings = load_settings() + + ws_url = build_ws_url() + headers = build_headers() + request = build_subscribe_request(SYMBOLS) + + print(f"WebSocket URL: {ws_url}") + print(f"Symbols: {', '.join(SYMBOLS)}") + print(f"Maximum messages: {MAX_MESSAGES}") + print(f"Total timeout: {TOTAL_TIMEOUT_SECONDS:.0f} seconds") + print() + + print("Subscription request:") + print( + json.dumps( + request, + ensure_ascii=False, + indent=2, + ) + ) + print() + + async with websockets.connect( + ws_url, + extra_headers=headers, + subprotocols=[Subprotocol("json")], + ping_interval=20, + ping_timeout=float(settings.exchange_timeout_sec), + open_timeout=float(settings.exchange_timeout_sec), + close_timeout=float(settings.exchange_timeout_sec), + ) as websocket: + await websocket.send( + json.dumps( + request, + ensure_ascii=False, + ) + ) + + print("Subscription request sent.") + print("Waiting for acknowledgement and trade events...") + print() + + loop = asyncio.get_running_loop() + started_at = loop.time() + deadline = started_at + TOTAL_TIMEOUT_SECONDS + + received_count = 0 + trade_event_count = 0 + timeout_count = 0 + + while received_count < MAX_MESSAGES: + remaining_seconds = deadline - loop.time() + + if remaining_seconds <= 0: + print( + "Total timeout reached after " + f"{elapsed_seconds(loop, started_at):.1f} seconds." + ) + break + + receive_timeout = min( + RECEIVE_TIMEOUT_SECONDS, + remaining_seconds, + ) + + try: + raw_message = await asyncio.wait_for( + websocket.recv(), + timeout=receive_timeout, + ) + except asyncio.TimeoutError: + timeout_count += 1 + + print( + "No trade event yet; connection is alive. " + f"Elapsed: {elapsed_seconds(loop, started_at):.1f}s, " + f"remaining: {max(0.0, remaining_seconds):.1f}s." + ) + continue + + if not isinstance( + raw_message, + (str, bytes, bytearray, memoryview), + ): + print( + "Skipped unsupported WebSocket message type: " + f"{type(raw_message).__name__}" + ) + continue + + received_count += 1 + document = decode_document(raw_message) + + if is_trade_event(document): + trade_event_count += 1 + message_kind = "TRADE EVENT" + else: + message_kind = "CONTROL / ACKNOWLEDGEMENT" + + print("=" * 80) + print( + f"Message #{received_count} — {message_kind} — " + f"elapsed {elapsed_seconds(loop, started_at):.1f}s" + ) + print(format_json(raw_message)) + print() + + print("=" * 80) + print("Diagnostic summary") + print(f"Received messages: {received_count}") + print(f"Trade events: {trade_event_count}") + print(f"Receive timeouts: {timeout_count}") + print( + "Elapsed time: " + f"{elapsed_seconds(loop, started_at):.1f} seconds" + ) + + +def main() -> None: + try: + asyncio.run(receive_trades()) + except KeyboardInterrupt: + print() + print("Stopped by user.") + except Exception as exc: + print() + print( + "WebSocket diagnostic failed: " + f"{type(exc).__name__}: {exc}" + ) + raise + + +if __name__ == "__main__": + main() \ No newline at end of file diff --git a/docs/migrations/build_057.md b/docs/migrations/build_057.md new file mode 100644 index 0000000..7b65724 --- /dev/null +++ b/docs/migrations/build_057.md @@ -0,0 +1,1111 @@ +# Build 057 — Исследование REST/WebSocket Trade Feed (Time & Sales) + +**Engineering Migration Report** + +--- + +# Контроль документа + +| Свойство | Значение | +|----------|----------| +| Build | 057 | +| Название | Исследование REST/WebSocket Trade Feed (Time & Sales) | +| Статус | Completed | +| Проект | Dzentra | +| Подсистема | Market Data Acquisition | +| Компонент | Trades Feed | +| Версия | 1.0 | +| Дата | 2026-07-16 | + +--- + +# Цель Build + +После завершения полной миграции OHLCV было принято решение перейти к следующему типу рыночных данных — **Trade Feed (Time & Sales)**. + +Перед реализацией нового канонического модуля требовалось определить: + +- существует ли полноценный REST API; +- существует ли полноценный WebSocket Feed; +- совпадают ли данные REST и WebSocket; +- какие поля действительно предоставляет биржа; +- имеются ли скрытые ограничения; +- какие особенности необходимо учитывать при проектировании новой архитектуры. + +Главной целью Build являлось **не написание кода**, а инженерное исследование поведения API. + +--- + +# Исследуемые интерфейсы + +Были исследованы два независимых источника данных. + +## REST + +``` +GET /api/v2/aggTrades +``` + +Назначение: + +получение последних агрегированных сделок. + +Swagger предоставляет параметры: + +- symbol +- fromId +- startTime +- endTime +- limit + +--- + +## WebSocket + +``` +destination = trades.subscribe +``` + +назначение: + +подписка на поток новых сделок. + +Получаемые события: + +``` +destination = internal.trade +``` + +--- + +# Выполненные задачи + +В рамках Build были последовательно проверены: + +- REST endpoint; +- WebSocket endpoint; +- формат ответа; +- соответствие Swagger документации; +- работа Demo; +- работа Production; +- соответствие REST и WebSocket; +- стабильность подписки; +- возможность дальнейшего использования в Dzentra. + +--- + +# Проверка REST API + +Для исследования использовались реальные запросы. + +Основная команда: + +```bash +curl -sS -G \ + "https://api-adapter.dzengi.com/api/v2/aggTrades" \ + --data-urlencode "symbol=BTC/USD_LEVERAGE" \ + --data-urlencode "limit=10" \ +| python -m json.tool +``` + +--- + +# Полученный формат REST + +Каждая запись имеет вид: + +```json +{ + "a": 2134857062, + "p": "64497.25", + "q": "0.005", + "T": 1784218066823, + "m": false +} +``` + +--- + +# Расшифровка полей + +| Поле | Значение | +|------|----------| +| a | идентификатор сделки | +| p | цена исполнения | +| q | объём сделки | +| T | timestamp сделки | +| m | сторона инициатора сделки | + +--- + +# Особенности REST + +Во время исследования было установлено: + +- данные отсортированы по времени; +- записи идут от новых к старым; +- timestamp имеет миллисекундную точность; +- объём может возвращаться в научной записи (`1.0E-4`); +- API поддерживает выборку по limit; +- идентификатор сделки монотонно возрастает. + +--- + +# Проверка Demo окружения + +Первоначально исследования проводились через Demo API. + +Использовался endpoint: + +```text +https://demo-api-adapter.dzengi.com/api/v2/aggTrades +``` + +Во время длительного наблюдения было обнаружено: + +- результаты практически не меняются; +- новые сделки появляются крайне редко; +- иногда ответы полностью совпадают на протяжении нескольких минут. + +Из этого был сделан вывод, что Demo нельзя использовать для проверки непрерывного потока сделок. + +--- + +# Проверка Production + +После перехода на Production API ситуация изменилась. + +Использовался endpoint: + +```text +https://api-adapter.dzengi.com/api/v2/aggTrades +``` + +Повторные запросы начали возвращать постоянно обновляющийся поток последних сделок. + +Это подтвердило: + +- REST endpoint является рабочим; +- данные обновляются в реальном времени; +- Production подходит для проверки Time & Sales. + +--- + +# Проверка поддержки символов + +Во время тестирования было установлено: + +``` +BTC/USD +``` + +не поддерживается. + +REST возвращает ошибку: + +```json +{ + "code": -1128, + "msg": "Invalid symbol: BTC/USD" +} +``` + +При этом корректным является символ: + +``` +BTC/USD_LEVERAGE +``` + +Именно его необходимо использовать в дальнейшем при работе с Trade Feed. + +--- + +# Предварительные выводы по REST + +На данном этапе исследования было подтверждено: + +✔ endpoint полностью работоспособен; + +✔ формат данных соответствует Swagger; + +✔ поля достаточны для построения Time & Sales; + +✔ данные поступают в режиме реального времени; + +✔ Demo окружение не подходит для полноценного тестирования активности рынка; + +✔ дальнейшие исследования следует проводить исключительно через Production API. + +--- + +# Исследование WebSocket + +После завершения анализа REST был исследован поток WebSocket. + +Использовался раздел Swagger: + +``` +trades.subscribe +``` + +Документация описывает механизм подписки на поток сделок. + +--- + +# Формат подписки + +Клиент отправляет сообщение: + +```json +{ + "correlationId": "", + "destination": "trades.subscribe", + "payload": { + "symbols": [ + "BTC/USD_LEVERAGE" + ] + } +} +``` + +После успешной обработки сервер отвечает подтверждением. + +--- + +# ACK-сообщение + +Первое сообщение всегда имеет вид: + +```json +{ + "correlationId": "...", + "destination": "trades.subscribe", + "payload": { + "subscriptions": { + "BTC/USD_LEVERAGE": "PROCESSED" + } + }, + "status": "OK" +} +``` + +Это сообщение не содержит рыночных данных. + +Оно означает исключительно успешную регистрацию подписки. + +--- + +# Проверка нескольких символов + +Для проверки обработки ошибок была отправлена подписка одновременно на два символа: + +```json +{ + "symbols": [ + "BTC/USD", + "BTC/USD_LEVERAGE" + ] +} +``` + +Ответ сервера: + +```json +{ + "subscriptions": { + "BTC/USD": "ERROR: BTC/USD not found", + "BTC/USD_LEVERAGE": "PROCESSED" + } +} +``` + +Из этого был сделан важный вывод. + +Подписка обрабатывается отдельно для каждого инструмента. + +Ошибка одного символа не приводит к отмене остальных подписок. + +--- + +# Первоначальное тестирование Demo + +Первые проверки выполнялись на Demo окружении. + +Подписка проходила успешно. + +Получалось корректное ACK. + +Однако далее: + +- сообщения сделок не приходили; +- соединение оставалось открытым; +- Ping/Pong работал корректно; +- сервер не разрывал соединение. + +Несмотря на пятиминутное ожидание поток оставался пустым. + +--- + +# Вывод по Demo + +Полученные результаты полностью совпали с поведением REST. + +Demo практически не генерирует торговую активность. + +Поэтому отсутствие сообщений не является ошибкой клиента. + +--- + +# Переход на Production + +Для окончательной проверки было принято решение перейти на Production API. + +Использовались настройки: + +```text +EXCHANGE_BASE_URL=https://api-adapter.dzengi.com +EXCHANGE_WS_URL=wss://api-adapter.dzengi.com +``` + +После этого поток начал немедленно передавать сделки. + +--- + +# Первый Trade Event + +Первое сообщение выглядело следующим образом. + +```json +{ + "destination": "internal.trade", + "payload": { + "buyer": false, + "id": 2134846831, + "orderId": "00a02503-0079-54c4-0000-000081e62b58", + "price": 64555.55, + "size": 0.002, + "symbol": "BTC/USD_LEVERAGE", + "ts": 1784218012030 + }, + "status": "OK" +} +``` + +Именно это сообщение является каноническим Trade Event. + +--- + +# Структура Trade Event + +Сервер публикует сообщения типа + +``` +destination = internal.trade +``` + +Поля сообщения: + +| Поле | Назначение | +|------|------------| +| id | уникальный идентификатор сделки | +| price | цена исполнения | +| size | объём | +| ts | время исполнения | +| buyer | направление сделки | +| symbol | торговый инструмент | +| orderId | внутренний идентификатор биржи | + +--- + +# Семантика buyer + +Экспериментально подтверждено: + +``` +buyer = true +``` + +означает + +> инициатор сделки — покупатель. + +Соответственно + +``` +buyer = false +``` + +означает + +> инициатор сделки — продавец. + +Данное поле полностью соответствует REST-полю: + +``` +m +``` + +по следующему правилу: + +``` +buyer == !m +``` + +То есть: + +| REST m | WebSocket buyer | +|---------|-----------------| +| true | false | +| false | true | + +Это полностью соответствует классической Binance-семантике: + +``` +m = isBuyerMaker +``` + +--- + +# Частота сообщений + +Во время наблюдения было получено множество последовательных сообщений. + +Интервал между ними не фиксирован. + +Были зарегистрированы: + +- менее секунды; +- несколько секунд; +- десятки секунд. + +Таким образом поток является событийным. + +Сообщения появляются только после фактического исполнения новой сделки. + +--- + +# Стабильность соединения + +Во время длительного тестирования: + +- соединение не разрывалось; +- Ping/Pong работал корректно; +- подписка сохранялась; +- повторной регистрации не требовалось. + +Это подтверждает возможность использования данного WebSocket для постоянной фоновой работы Dzentra. + +--- + +# Диагностический скрипт + +Для исследования был разработан отдельный диагностический инструмент: + +``` +scripts/check_trades_websocket.py +``` + +Назначение: + +- открытие WebSocket; +- регистрация подписки; +- вывод ACK; +- вывод всех Trade Event; +- отображение времени получения сообщений; +- автоматический Ping; +- контроль таймаутов; +- длительное наблюдение за потоком. + +Скрипт не является частью рабочей архитектуры. + +Он предназначен исключительно для инженерной диагностики и проверки API. + +--- + +# Запуск скрипта + +Для Demo: + +```bash +PYTHONPATH=. python scripts/check_trades_websocket.py +``` + +Для Production: + +```bash +EXCHANGE_BASE_URL=https://api-adapter.dzengi.com \ +EXCHANGE_WS_URL=wss://api-adapter.dzengi.com \ +EXCHANGE_API_KEY="" \ +PYTHONPATH=. \ +python scripts/check_trades_websocket.py +``` + +--- + +# Проверенные сценарии + +Во время исследования были успешно проверены: + +- успешная регистрация подписки; +- получение ACK; +- получение Trade Event; +- получение нескольких подряд Trade Event; +- работа Ping/Pong; +- длительное соединение; +- отсутствие ложных сообщений; +- корректная обработка нескольких символов; +- обработка ошибки неизвестного символа. + +--- + +# Сопоставление REST и WebSocket + +Главной задачей исследования было доказать, что REST и WebSocket описывают **одни и те же биржевые сделки**, а не различные представления рынка. + +Для этого была выполнена серия сравнительных экспериментов. + +--- + +# Методика проверки + +Алгоритм проверки был следующим. + +1. Запустить подписку `trades.subscribe`. +2. Дождаться получения новой сделки через WebSocket. +3. Зафиксировать: + - id; + - цену; + - объём; + - timestamp; + - направление сделки. +4. Немедленно запросить последние сделки через REST. +5. Найти сделку по идентификатору. + +--- + +# Проверка №1 + +Через WebSocket была получена сделка: + +```json +{ + "buyer": false, + "id": 2134846831, + "price": 64555.55, + "size": 0.002, + "symbol": "BTC/USD_LEVERAGE", + "ts": 1784218012030 +} +``` + +После этого был выполнен REST-запрос: + +```bash +curl -G \ +https://api-adapter.dzengi.com/api/v2/aggTrades \ +--data-urlencode "symbol=BTC/USD_LEVERAGE" \ +--data-urlencode "limit=100" +``` + +Далее был выполнен поиск по ID сделки. + +Полученный результат: + +```json +{ + "a": 2134846831, + "p": "64555.55", + "q": "0.002", + "T": 1784218012030, + "m": true +} +``` + +--- + +# Совпадение полей + +Было подтверждено полное совпадение. + +| REST | WebSocket | +|-------|-----------| +| a | id | +| p | price | +| q | size | +| T | ts | +| m | !buyer | + +Все значения совпали полностью. + +--- + +# Проверка нескольких сообщений + +Аналогичная проверка была выполнена ещё для нескольких последовательных событий. + +Каждый раз наблюдалось одинаковое соответствие: + +- идентификатор сделки совпадает; +- цена совпадает; +- объём совпадает; +- timestamp совпадает; +- направление определяется инверсией поля `buyer`. + +Таким образом экспериментально доказано, что REST и WebSocket используют единую торговую базу данных. + +--- + +# Итоговая схема соответствия + +```text + Exchange Matching Engine + │ + ┌─────────────┴─────────────┐ + │ │ + REST aggTrades WebSocket trades.subscribe + │ │ + ▼ ▼ + Aggregated Trade internal.trade + │ │ + └─────────────┬─────────────┘ + ▼ + Canonical Trade Model +``` + +Это означает, что в архитектуре Dzentra оба источника должны преобразовываться **в одну каноническую модель**, а не существовать как независимые типы данных. + +--- + +# Архитектурные выводы + +По результатам исследования были сформированы следующие инженерные выводы. + +## 1. REST и WebSocket описывают одну сущность + +Различается только транспорт доставки. + +REST используется для получения истории. + +WebSocket используется для получения новых событий. + +Модель данных должна быть общей. + +--- + +## 2. Канонической сущностью становится Trade + +Не существует причин создавать отдельные модели: + +- RestTrade; +- WsTrade; +- InternalTrade. + +Должна существовать одна модель: + +```text +Trade +``` + +которая используется всеми компонентами системы. + +--- + +## 3. Trade Feed должен быть независимым модулем + +Как и Candle Feed, новый модуль должен располагаться в подсистеме Market Data Acquisition. + +Предварительная структура: + +```text +market_data/acquisition/ + + models/ + trade.py + + feeds/ + trades_feed.py + + handlers/ + trades_handler.py + + adapters/dzengi/ + rest.py + websocket.py + + validation/ + schema.py + values.py +``` + +--- + +## 4. REST и WebSocket должны использовать общий Handler + +Различаться должны только источники получения документов. + +Схема обработки должна повторять архитектуру Candle Feed: + +```text +REST Document + │ + ▼ +Parser + │ + ▼ +Validation + │ + ▼ +Mapper + │ + ▼ +Trade +``` + +и + +```text +WebSocket Document + │ + ▼ +Parser + │ + ▼ +Validation + │ + ▼ +Mapper + │ + ▼ +Trade +``` + +--- + +## 5. Внутри Dzentra не должно существовать транспортных моделей + +После обработки все компоненты должны работать исключительно с: + +```text +Trade +``` + +Они не должны знать: + +- откуда пришли данные; +- были ли они получены через REST; +- были ли они получены через WebSocket. + +--- + +# Возможности будущего Trade Feed + +После реализации модуль сможет использоваться для решения задач, которые невозможно решить только по свечам. + +В частности: + +- построение Time & Sales; +- определение текущего торгового потока; +- оценка направления агрессора (Buyer / Seller); +- обнаружение крупных рыночных сделок; +- поиск всплесков активности; +- вычисление Buy/Sell Imbalance; +- вычисление Trade Delta; +- обнаружение поглощений; +- обнаружение кульминаций объёма; +- анализ последовательности исполнения сделок; +- подтверждение истинности пробоев; +- подтверждение ложных пробоев; +- анализ микроструктуры рынка. + +Все эти алгоритмы будут использовать одну и ту же каноническую модель `Trade`. + +--- + +# Итоги Build 057 + +В рамках Build были полностью исследованы REST и WebSocket интерфейсы Trade Feed. + +Подтверждено: + +- ✔ REST `GET /api/v2/aggTrades` полностью работоспособен. +- ✔ WebSocket `trades.subscribe` полностью работоспособен. +- ✔ Demo окружение практически не генерирует поток сделок и не подходит для тестирования. +- ✔ Production окружение генерирует полноценный поток Trade Event. +- ✔ REST и WebSocket описывают одну и ту же сделку. +- ✔ Идентификаторы, цены, объёмы и временные метки полностью совпадают. +- ✔ Поле `buyer` является логической инверсией REST-поля `m`. +- ✔ Разработан и протестирован диагностический скрипт `scripts/check_trades_websocket.py`. +- ✔ Подготовлена архитектурная база для реализации канонического Trade Feed. + +--- + +# Следующий этап + +Следующим Build станет **Build 058 — Canonical Trade Model Foundation**. + +В его рамках будет разработан первый полноценный компонент нового Trade Feed: + +- каноническая модель `Trade`; +- REST Parser; +- REST Validation; +- REST Mapper; +- базовые unit-тесты; +- интеграция в существующую архитектуру `Market Data Acquisition`. + +После этого система Dzentra получит вторую полноценную каноническую рыночную сущность после завершённой миграции OHLCV — поток реальных биржевых сделок (Time & Sales). + +# Приложение A — Использование диагностического скрипта Trade WebSocket + +--- + +# Назначение скрипта + +В рамках Build 057 был разработан отдельный диагностический инструмент + +``` +scripts/check_trades_websocket.py +``` + +Скрипт предназначен исключительно для инженерной проверки работы биржевого +Trade WebSocket API. + +Он не является частью рабочего Runtime Dzentra и не используется торговой +логикой. + +Его задачи: + +- проверка открытия WebSocket-соединения; +- проверка успешной регистрации подписки; +- проверка получения ACK; +- проверка получения новых сделок; +- проверка корректности Ping/Pong; +- длительное наблюдение за потоком; +- диагностика работы Demo и Production окружений; +- подтверждение соответствия WebSocket и REST. + +--- + +# Проверяемые возможности + +Во время работы скрипт автоматически проверяет: + +- возможность подключения к серверу; +- правильность формирования запроса подписки; +- успешность регистрации подписки; +- получение сообщений подтверждения; +- получение сообщений Trade Event; +- сохранение соединения при отсутствии новых сделок; +- работу Ping/Pong; +- получение нескольких последовательных сделок; +- возможность длительной непрерывной работы. + +--- + +# Используемые переменные окружения + +Скрипт использует стандартные настройки проекта. + +Основные переменные: + +``` +EXCHANGE_BASE_URL +``` + +``` +EXCHANGE_WS_URL +``` + +``` +EXCHANGE_API_KEY +``` + +Если переменные не заданы, используются значения из конфигурации Dzentra. + +--- + +# Запуск Demo + +```bash +PYTHONPATH=. python scripts/check_trades_websocket.py +``` + +Используется: + +``` +wss://demo-api-adapter.dzengi.com/connect +``` + +--- + +# Ожидаемый результат Demo + +При корректной работе вывод должен быть примерно следующим. + +``` +Subscription request sent. + +Message #1 + +ACK +``` + +После этого возможно длительное отсутствие сообщений Trade Event. + +Периодически выводится сообщение: + +``` +No trade event yet; connection is alive. +``` + +Это означает: + +- соединение активно; +- Ping/Pong работает; +- новых сделок пока нет. + +Подобное поведение является нормальным для Demo окружения. + +--- + +# Запуск Production + +Для проверки реального рынка используется: + +```bash +EXCHANGE_BASE_URL=https://api-adapter.dzengi.com \ +EXCHANGE_WS_URL=wss://api-adapter.dzengi.com \ +EXCHANGE_API_KEY="" \ +PYTHONPATH=. \ +python scripts/check_trades_websocket.py +``` + +--- + +# Ожидаемый результат Production + +После ACK начинают поступать реальные сделки. + +Типичный вывод: + +``` +Message #2 + +TRADE EVENT +``` + +```json +{ + "destination": "internal.trade", + "payload": { + "id": 2134846831, + "price": 64555.55, + "size": 0.002, + "buyer": false, + "symbol": "BTC/USD_LEVERAGE", + "ts": 1784218012030 + } +} +``` + +Затем по мере появления новых сделок будут выводиться: + +``` +Message #3 + +TRADE EVENT +``` + +``` +Message #4 + +TRADE EVENT +``` + +и так далее. + +--- + +# Интерпретация сообщений + +Во время работы возможны два типа сообщений. + +## ACK + +```text +destination = trades.subscribe +``` + +означает успешную регистрацию подписки. + +--- + +## TRADE EVENT + +```text +destination = internal.trade +``` + +означает получение новой сделки с биржи. + +--- + +# Проверка соответствия REST + +После получения Trade Event рекомендуется выполнить запрос: + +```bash +curl -G \ +"https://api-adapter.dzengi.com/api/v2/aggTrades" \ +--data-urlencode "symbol=BTC/USD_LEVERAGE" \ +--data-urlencode "limit=100" +``` + +Затем убедиться, что сделка с полученным ID присутствует в REST. + +Пример поиска: + +```bash +curl -sS -G \ +"https://api-adapter.dzengi.com/api/v2/aggTrades" \ +--data-urlencode "symbol=BTC/USD_LEVERAGE" \ +--data-urlencode "limit=100" \ +| python -c ' + +import json +import sys + +target_id = 2134846831 + +items = json.load(sys.stdin) + +matches = [ + item + for item in items + if item["a"] == target_id +] + +print(json.dumps(matches, indent=2)) +' +``` + +Если вывод содержит одну запись с тем же идентификатором, значит соответствие REST и WebSocket подтверждено. + +--- + +# Диагностические признаки корректной работы + +Работа скрипта считается полностью успешной, если выполняются все условия: + +- успешно устанавливается WebSocket-соединение; +- получено ACK; +- подписка имеет статус `PROCESSED`; +- поступают сообщения `internal.trade`; +- соединение не разрывается при длительном ожидании; +- Ping/Pong выполняется успешно; +- сделки из WebSocket присутствуют в REST `aggTrades`. + +--- + +# Итог + +Диагностический скрипт полностью подтвердил работоспособность биржевого интерфейса Trade Feed и стал основным инженерным инструментом проверки перед началом реализации канонического модуля **Trade (Time & Sales)** в Build 058. \ No newline at end of file