build 059.6: map websocket OHLC close events

This commit is contained in:
2026-07-17 09:59:11 +03:00
parent 5c7ecf9340
commit 85e9662f5b
4 changed files with 1158 additions and 0 deletions

View File

@@ -14,19 +14,23 @@ from src.market_data.acquisition.adapters.dzengi.models import (
DzengiMinNotionalFilter,
DzengiRawNumeric,
DzengiTicker24hrResponse,
DzengiWebSocketOhlcEvent,
DzengiWebSocketQuoteResponse,
)
from src.market_data.acquisition.exceptions import (
CandleMappingError,
CandleWebSocketMappingError,
InstrumentReferenceMappingError,
QuoteMappingError,
)
from src.market_data.acquisition.models.candle import Candle
from src.market_data.acquisition.models.candle_close import CandleCloseEvent
from src.market_data.acquisition.models.instrument import Instrument
from src.market_data.acquisition.models.quote import Quote
_DZENGI_SOURCE_NAME = "dzengi"
_DZENGI_WEBSOCKET_OHLC_SOURCE_NAME = "dzengi_websocket_ohlc"
def map_dzengi_symbol_to_instrument(
@@ -318,6 +322,104 @@ def map_dzengi_websocket_quote_to_quote(
source=_DZENGI_SOURCE_NAME,
)
def map_dzengi_websocket_ohlc_to_candle_close_event(
event: DzengiWebSocketOhlcEvent,
*,
received_at: datetime,
) -> CandleCloseEvent:
"""
Преобразовать проверенное Dzengi WebSocket OHLC-событие
во внутреннюю модель CandleCloseEvent.
Функция предполагает, что до mapper уже были выполнены:
schema validation, parsing и value validation.
"""
normalized_received_at = _require_websocket_ohlc_aware_datetime(
received_at,
field_name="received_at",
)
return CandleCloseEvent(
symbol=event.symbol.strip(),
interval=event.interval,
candle_type=event.candle_type,
open_time=_websocket_ohlc_timestamp_ms_to_utc_datetime(
event.open_time,
),
received_at=normalized_received_at,
open_price=_required_websocket_ohlc_decimal(
event.open_price,
field_name="open",
),
high_price=_required_websocket_ohlc_decimal(
event.high_price,
field_name="high",
),
low_price=_required_websocket_ohlc_decimal(
event.low_price,
field_name="low",
),
close_price=_required_websocket_ohlc_decimal(
event.close_price,
field_name="close",
),
source=_DZENGI_WEBSOCKET_OHLC_SOURCE_NAME,
)
def _websocket_ohlc_timestamp_ms_to_utc_datetime(
value: int,
) -> datetime:
try:
return datetime.fromtimestamp(
value / 1000,
tz=timezone.utc,
)
except (OverflowError, OSError, ValueError) as exc:
raise CandleWebSocketMappingError(
"Поле t WebSocket OHLC-события невозможно "
"преобразовать в UTC datetime."
) from exc
def _required_websocket_ohlc_decimal(
value: DzengiRawNumeric,
*,
field_name: str,
) -> Decimal:
try:
result = Decimal(str(value))
except (InvalidOperation, ValueError) as exc:
raise CandleWebSocketMappingError(
f"Поле {field_name} WebSocket OHLC-события невозможно "
"преобразовать в Decimal."
) from exc
if not result.is_finite():
raise CandleWebSocketMappingError(
f"Поле {field_name} WebSocket OHLC-события должно быть "
"конечным числом."
)
return result
def _require_websocket_ohlc_aware_datetime(
value: datetime,
*,
field_name: str,
) -> datetime:
if value.tzinfo is None or value.utcoffset() is None:
raise CandleWebSocketMappingError(
f"Поле {field_name} должно содержать "
"timezone-aware datetime."
)
return value
def map_dzengi_klines_to_candles(
response: DzengiKlinesResponse,
*,

View File

@@ -93,6 +93,12 @@ class CandleWebSocketValueError(MarketDataAcquisitionError):
pass
# Ошибка преобразования WebSocket OHLC-события
# во внутреннюю модель уведомления о закрытии свечи.
class CandleWebSocketMappingError(MarketDataAcquisitionError):
pass
# Ошибка преобразования проверенного документа в raw-модели свечей.
class CandleParseError(MarketDataAcquisitionError):
pass

View File

@@ -0,0 +1,253 @@
# app/tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_mapper.py
from __future__ import annotations
from datetime import datetime, timedelta, timezone
from decimal import Decimal
from typing import TypedDict
import pytest
from src.market_data.acquisition.adapters.dzengi.mapper import (
map_dzengi_websocket_ohlc_to_candle_close_event,
)
from src.market_data.acquisition.adapters.dzengi.models import (
DzengiWebSocketOhlcEvent,
)
from src.market_data.acquisition.exceptions import (
CandleWebSocketMappingError,
)
from src.market_data.acquisition.models.candle_close import (
CandleCloseEvent,
)
class OhlcValues(TypedDict):
open_price: str | int | float
high_price: str | int | float
low_price: str | int | float
close_price: str | int | float
def _event(
*,
open_time: int = 1784227140000,
open_price: str | int | float = "63992.00",
high_price: str | int | float = "64032.55",
low_price: str | int | float = "63984.00",
close_price: str | int | float = "64032.55",
) -> DzengiWebSocketOhlcEvent:
return DzengiWebSocketOhlcEvent(
symbol=" BTC/USD_LEVERAGE ",
interval="1m",
candle_type="classic",
open_time=open_time,
open_price=open_price,
high_price=high_price,
low_price=low_price,
close_price=close_price,
)
def _received_at() -> datetime:
return datetime(
2026,
7,
16,
18,
40,
0,
125000,
tzinfo=timezone.utc,
)
def test_map_websocket_ohlc_returns_candle_close_event() -> None:
result = map_dzengi_websocket_ohlc_to_candle_close_event(
_event(),
received_at=_received_at(),
)
assert isinstance(result, CandleCloseEvent)
def test_map_websocket_ohlc_maps_all_fields() -> None:
result = map_dzengi_websocket_ohlc_to_candle_close_event(
_event(),
received_at=_received_at(),
)
assert result.symbol == "BTC/USD_LEVERAGE"
assert result.interval == "1m"
assert result.candle_type == "classic"
assert result.open_time == datetime(
2026,
7,
16,
18,
39,
tzinfo=timezone.utc,
)
assert result.received_at == _received_at()
assert result.open_price == Decimal("63992.00")
assert result.high_price == Decimal("64032.55")
assert result.low_price == Decimal("63984.00")
assert result.close_price == Decimal("64032.55")
assert result.source == "dzengi_websocket_ohlc"
def test_map_websocket_ohlc_converts_numeric_variants_to_decimal() -> None:
result = map_dzengi_websocket_ohlc_to_candle_close_event(
_event(
open_price=63992,
high_price=64032.55,
low_price="63984.00",
close_price="64001.25",
),
received_at=_received_at(),
)
assert result.open_price == Decimal("63992")
assert result.high_price == Decimal("64032.55")
assert result.low_price == Decimal("63984.00")
assert result.close_price == Decimal("64001.25")
def test_map_websocket_ohlc_preserves_aware_received_at() -> None:
received_at = datetime(
2026,
7,
16,
21,
40,
tzinfo=timezone(timedelta(hours=3)),
)
result = map_dzengi_websocket_ohlc_to_candle_close_event(
_event(),
received_at=received_at,
)
assert result.received_at is received_at
assert result.received_at.utcoffset() == timedelta(hours=3)
def test_map_websocket_ohlc_rejects_naive_received_at() -> None:
received_at = datetime(
2026,
7,
16,
18,
40,
)
with pytest.raises(
CandleWebSocketMappingError,
match="received_at.*timezone-aware",
):
map_dzengi_websocket_ohlc_to_candle_close_event(
_event(),
received_at=received_at,
)
def test_map_websocket_ohlc_rejects_invalid_timestamp() -> None:
with pytest.raises(
CandleWebSocketMappingError,
match="Поле t.*UTC datetime",
):
map_dzengi_websocket_ohlc_to_candle_close_event(
_event(open_time=10**30),
received_at=_received_at(),
)
@pytest.mark.parametrize(
("field_name", "field_value"),
(
("open_price", "invalid"),
("high_price", "invalid"),
("low_price", "invalid"),
("close_price", "invalid"),
),
)
def test_map_websocket_ohlc_rejects_invalid_decimal(
field_name: str,
field_value: str,
) -> None:
values: OhlcValues = {
"open_price": "63992.00",
"high_price": "64032.55",
"low_price": "63984.00",
"close_price": "64032.55",
}
values[field_name] = field_value
with pytest.raises(CandleWebSocketMappingError):
map_dzengi_websocket_ohlc_to_candle_close_event(
_event(**values),
received_at=_received_at(),
)
@pytest.mark.parametrize(
("field_name", "field_value"),
(
("open_price", "NaN"),
("high_price", "Infinity"),
("low_price", "-Infinity"),
("close_price", "NaN"),
),
)
def test_map_websocket_ohlc_rejects_non_finite_decimal(
field_name: str,
field_value: str,
) -> None:
values: OhlcValues = {
"open_price": "63992.00",
"high_price": "64032.55",
"low_price": "63984.00",
"close_price": "64032.55",
}
values[field_name] = field_value
with pytest.raises(
CandleWebSocketMappingError,
match="конечным числом",
):
map_dzengi_websocket_ohlc_to_candle_close_event(
_event(**values),
received_at=_received_at(),
)
def test_map_websocket_ohlc_preserves_heikin_ashi_type() -> None:
event = DzengiWebSocketOhlcEvent(
symbol="BTC/USD_LEVERAGE",
interval="1m",
candle_type="heikin-ashi",
open_time=1784227140000,
open_price="64068.18",
high_price="64128.80",
low_price="64068.18",
close_price="64100.85",
)
result = map_dzengi_websocket_ohlc_to_candle_close_event(
event,
received_at=_received_at(),
)
assert result.candle_type == "heikin-ashi"
def test_map_websocket_ohlc_does_not_create_canonical_volume() -> None:
result = map_dzengi_websocket_ohlc_to_candle_close_event(
_event(),
received_at=_received_at(),
)
assert not hasattr(result, "volume")

View File

@@ -0,0 +1,797 @@
# Build 059.6 — WebSocket OHLC Mapper
**Проект:** Dzentra
**Подсистема:** Market Data Acquisition
**Этап:** 059.6
**Статус:** Completed
---
# Цель
Реализовать преобразование проверенного WebSocket OHLC-события Dzengi во внутреннюю модель Dzentra:
```text
CandleCloseEvent
```
Mapper должен стать последней стадией обработки WebSocket-сообщения перед передачей события во внутренний runtime.
---
# Причина изменения
После завершения Build 059.5 система уже содержит:
```text
Raw WebSocket JSON
Schema Validation
ValidatedWebSocketOhlcDocument
Parser
DzengiWebSocketOhlcEvent
Value Validation
```
Однако transport-модель адаптера:
```text
DzengiWebSocketOhlcEvent
```
не должна использоваться за пределами слоя интеграции с Dzengi.
Для внутренних компонентов требуется полностью source-independent модель:
```text
CandleCloseEvent
```
Build 059.6 добавляет именно этот переход.
---
# Реализовано
Изменены файлы:
```text
src/market_data/acquisition/adapters/dzengi/mapper.py
src/market_data/acquisition/exceptions.py
```
Создан файл:
```text
tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_mapper.py
```
---
# Новый mapper
Добавлена функция:
```text
map_dzengi_websocket_ohlc_to_candle_close_event()
```
Назначение функции:
```text
DzengiWebSocketOhlcEvent
CandleCloseEvent
```
Mapper предполагает, что ранее уже были выполнены:
- schema validation;
- parsing;
- value validation.
Повторная проверка предметной корректности данных не выполняется.
---
# Выполняемые преобразования
Mapper выполняет только преобразования типов.
## Timestamp
Поле:
```text
t
```
преобразуется из:
```text
int (milliseconds)
```
в:
```text
datetime (UTC)
```
---
## OHLC
Поля:
```text
o
h
l
c
```
преобразуются из:
```text
str | int | float
```
в:
```text
Decimal
```
Все вычисления внутри runtime далее выполняются только с использованием `Decimal`.
---
## received_at
Mapper принимает дополнительный аргумент:
```text
received_at
```
Тип:
```text
datetime
```
Обязательное требование:
```text
timezone-aware datetime
```
Naive datetime считается ошибкой mapper.
---
## source
Mapper автоматически устанавливает источник:
```text
dzengi_websocket_ohlc
```
Тем самым внутренние модели больше не зависят от транспортного контракта адаптера.
---
# Новый тип ошибки
Добавлено исключение:
```text
CandleWebSocketMappingError
```
Исключение используется исключительно на этапе mapping.
Это позволяет разделить ошибки различных стадий обработки.
Схема обработки стала следующей:
```text
Transport
CandleTransportError
Schema
CandleWebSocketSchemaError
Parser
CandleWebSocketParseError
Value Validation
CandleWebSocketValueError
Mapper
CandleWebSocketMappingError
```
Таким образом каждая стадия имеет собственный независимый контракт ошибок.
---
# Почему не используется CandleMappingError
В проекте уже существует:
```text
CandleMappingError
```
Однако он относится исключительно к REST Candles Feed.
REST mapper строит:
```text
Canonical Candle
```
WebSocket mapper строит:
```text
CandleCloseEvent
```
Это разные сущности.
Использование общего исключения привело бы к смешиванию двух различных контрактов.
Поэтому был введён отдельный тип ошибки.
---
# Использование существующих helper-функций
В mapper уже присутствуют helper-функции для REST-свечей.
Например:
```text
_required_candle_decimal()
```
Они намеренно не переиспользуются.
Причина:
они выбрасывают:
```text
CandleMappingError
```
WebSocket mapper обязан выбрасывать:
```text
CandleWebSocketMappingError
```
Поэтому были реализованы отдельные helper-функции.
---
# Новые helper-функции
Добавлены:
```text
_required_websocket_ohlc_decimal()
_websocket_ohlc_timestamp_ms_to_utc_datetime()
_require_websocket_ohlc_aware_datetime()
```
Они полностью изолированы от REST mapper.
Это исключает смешивание различных источников данных и их контрактов.
---
# Преобразование времени
Поле:
```text
t
```
содержит Unix timestamp в миллисекундах.
Mapper преобразует его через:
```python
datetime.fromtimestamp(
value / 1000,
tz=timezone.utc,
)
```
В результате:
```text
open_time
```
становится объектом:
```text
datetime UTC
```
Внутренние слои больше не работают с миллисекундными timestamp.
---
# Преобразование Decimal
Каждое поле:
```text
open
high
low
close
```
проходит преобразование:
```text
str | int | float
Decimal
```
Дополнительно проверяется:
- корректность числа;
- конечность значения.
---
# Контракт mapper
Mapper не выполняет:
- schema validation;
- parsing;
- value validation;
- REST reconciliation;
- получение объёма;
- создание канонической `Candle`.
Единственная задача mapper:
```text
построить внутреннюю модель
CandleCloseEvent
из уже проверенного
DzengiWebSocketOhlcEvent
```
---
# Что mapper НЕ делает
Mapper намеренно не выполняет следующие действия.
## Не создаёт volume
Фактический runtime-контракт Dzengi содержит:
```text
symbol
interval
type
t
o
h
l
c
```
Поля:
```text
volume
```
не существует.
Поэтому mapper не должен:
- подставлять `0`;
- использовать `None`;
- вычислять объём самостоятельно.
Полный объём будет получен позднее через REST reconciliation.
---
## Не создаёт Candle
Mapper не строит:
```text
Candle
```
Причина проста.
Каноническая модель требует:
```text
volume
```
которого WebSocket не содержит.
Следовательно единственным корректным результатом mapper является:
```text
CandleCloseEvent
```
---
## Не выполняет REST reconciliation
Build 059.6 заканчивается на построении внутреннего события.
Дальнейший pipeline выглядит следующим образом:
```text
CandleCloseEvent
REST reconciliation
Canonical Candle
```
Сам reconciliation будет реализован отдельным этапом.
---
# Поддерживаемые типы свечей
Mapper сохраняет значение:
```text
candle_type
```
без изменения.
Поддерживаемые runtime-значения:
```text
classic
heikin-ashi
```
Mapper не интерпретирует содержимое свечи.
Он лишь переносит значение в внутреннюю модель.
---
# Обработка received_at
Поле:
```text
received_at
```
не вычисляется автоматически.
Оно передаётся извне.
Это решение позволяет:
- использовать единый момент получения сообщения;
- исключить влияние задержек внутри mapper;
- обеспечить единый источник времени для всей системы.
Mapper только проверяет:
```text
timezone-aware datetime
```
Naive datetime приводит к ошибке:
```text
CandleWebSocketMappingError
```
---
# Конвейер после Build 059.6
После завершения этапа полный pipeline выглядит следующим образом:
```text
Raw WebSocket message
Schema Validation
ValidatedWebSocketOhlcDocument
Parser
DzengiWebSocketOhlcEvent
Value Validation
Mapper
CandleCloseEvent
```
На этом Build 059.6 завершается.
Следующий этап добавит использование новой модели внутри runtime.
---
# Unit Tests
Создан файл:
```text
tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_mapper.py
```
Проверяются:
- построение `CandleCloseEvent`;
- преобразование всех полей;
- преобразование timestamp;
- преобразование числовых типов в `Decimal`;
- сохранение `received_at`;
- поддержка `heikin-ashi`;
- обрезка пробелов у symbol;
- проверка timezone-aware datetime;
- ошибка при невозможности преобразовать timestamp;
- ошибка при некорректном Decimal;
- ошибка при бесконечных значениях (`NaN`, `Infinity`);
- отсутствие создания `volume`.
Всего выполнено:
```text
16 tests
```
---
# Проверка синтаксиса
Выполнена команда:
```bash
python -m compileall \
src/market_data/acquisition/exceptions.py \
src/market_data/acquisition/adapters/dzengi/mapper.py \
tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_mapper.py
```
Результат:
```text
успешно
```
---
# Целевые тесты
Выполнена команда:
```bash
python -m pytest -q \
tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_mapper.py
```
Результат:
```text
16 passed
```
---
# Регрессия mapper
Выполнена команда:
```bash
python -m pytest -q \
tests/unit/market_data/acquisition/adapters/dzengi/test_mapper.py \
tests/unit/market_data/acquisition/adapters/dzengi/test_quote_mapper.py \
tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_quote_mapper.py \
tests/unit/market_data/acquisition/adapters/dzengi/test_candle_mapper.py \
tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_mapper.py
```
Результат:
```text
52 passed
```
---
# Дополнительная регрессия
Выполнена команда:
```bash
python -m pytest -q \
tests/unit/market_data/acquisition/models \
tests/unit/market_data/acquisition/validation/test_websocket_ohlc_values.py
```
Результат:
```text
126 passed
```
---
# Проверка форматирования
Выполнена команда:
```bash
git diff --check
```
Ошибок форматирования не обнаружено.
---
# Изменённые файлы
```text
src/market_data/acquisition/adapters/dzengi/mapper.py
src/market_data/acquisition/exceptions.py
```
Создан файл:
```text
tests/unit/market_data/acquisition/adapters/dzengi/test_websocket_ohlc_mapper.py
```
Документация:
```text
docs/migrations/build_059_6.md
```
---
# Совместимость
Build 059.6 полностью обратно совместим.
Не изменены:
- REST Candles Feed;
- REST Candle mapper;
- Quote mapper;
- Instrument mapper;
- ExchangeService;
- runtime;
- Trading;
- Market Analysis;
- каноническая модель `Candle`.
Новый mapper пока не используется существующим runtime и не влияет на работу действующего торгового бота.
---
# Архитектурное значение
Build 059.6 завершает формирование полного внутреннего конвейера преобразования WebSocket OHLC.
После этого этапа система имеет чёткое разделение уровней ответственности:
```text
Transport Layer
DzengiWebSocketOhlcEvent
Validation Layer
Mapper Layer
CandleCloseEvent
```
Каждый уровень отвечает только за собственную задачу:
| Уровень | Ответственность |
|---------|-----------------|
| Schema Validation | Проверка структуры входящего сообщения |
| Parser | Построение transport-модели |
| Value Validation | Проверка допустимости значений |
| Mapper | Преобразование transport-модели во внутреннюю модель |
| Runtime | Использование внутренней модели |
Такое разделение соответствует общей архитектуре Dzentra и исключает смешивание обязанностей между слоями.
---
# Ограничения этапа
Build 059.6 намеренно **не реализует**:
- подписку на WebSocket;
- публикацию события в runtime;
- REST reconciliation;
- получение объёма свечи;
- создание канонической `Candle`;
- проверку соответствия REST и WebSocket OHLC;
- восстановление после разрыва соединения;
- повторную синхронизацию истории свечей.
Все перечисленные задачи относятся к последующим этапам Build 059.
---
# Итог
Build 059.6 добавляет полноценный mapper:
```text
DzengiWebSocketOhlcEvent
CandleCloseEvent
```
В результате система получила:
- источник-независимую внутреннюю модель события закрытия свечи;
- корректное преобразование timestamp в `datetime (UTC)`;
- преобразование цен в `Decimal`;
- собственный контракт ошибок (`CandleWebSocketMappingError`);
- полную изоляцию WebSocket-контракта Dzengi от внутренних компонентов Dzentra;
- набор unit-тестов, подтверждающих корректность преобразования и обработки ошибок.
После Build 059.6 цепочка обработки WebSocket OHLC стала полностью завершённой вплоть до внутреннего события Dzentra.
---
# Следующий этап
```text
Build 059.7 — WebSocket Adapter
```
На этом этапе новый mapper будет встроен в адаптер получения рыночных данных, чтобы при поступлении подтверждённого WebSocket OHLC-сообщения формировался `CandleCloseEvent`, который затем сможет быть передан в механизм REST reconciliation.