build 044: establish canonical candles feed foundation

This commit is contained in:
2026-07-15 10:26:55 +03:00
parent a6325cf4cf
commit 47c75d68b1
22 changed files with 2937 additions and 316 deletions

View File

@@ -9,6 +9,7 @@ from src.market_data.acquisition.adapters.dzengi.models import (
DzengiExchangeInfoResponse,
DzengiExchangeInfoSymbol,
DzengiInstrumentFilter,
DzengiKlinesResponse,
DzengiLotSizeFilter,
DzengiMinNotionalFilter,
DzengiRawNumeric,
@@ -16,9 +17,11 @@ from src.market_data.acquisition.adapters.dzengi.models import (
DzengiWebSocketQuoteResponse,
)
from src.market_data.acquisition.exceptions import (
CandleMappingError,
InstrumentReferenceMappingError,
QuoteMappingError,
)
from src.market_data.acquisition.models.candle import Candle
from src.market_data.acquisition.models.instrument import Instrument
from src.market_data.acquisition.models.quote import Quote
@@ -314,3 +317,91 @@ def map_dzengi_websocket_quote_to_quote(
received_at=normalized_received_at,
source=_DZENGI_SOURCE_NAME,
)
def map_dzengi_klines_to_candles(
response: DzengiKlinesResponse,
*,
symbol: str,
interval: str,
source: str,
) -> tuple[Candle, ...]:
"""
Преобразовать проверенные raw-модели Dzengi klines
в канонический immutable-набор Candle.
Функция предполагает, что до mapper уже были выполнены:
schema validation, parsing и value validation.
"""
candles = tuple(
Candle(
symbol=symbol.strip(),
interval=interval.strip(),
open_time=_candle_timestamp_ms_to_utc_datetime(
item.open_time,
),
open_price=_required_candle_decimal(
item.open_price,
field_name="open",
),
high_price=_required_candle_decimal(
item.high_price,
field_name="high",
),
low_price=_required_candle_decimal(
item.low_price,
field_name="low",
),
close_price=_required_candle_decimal(
item.close_price,
field_name="close",
),
volume=_required_candle_decimal(
item.volume,
field_name="volume",
),
source=source.strip(),
)
for item in response.items
)
return tuple(
sorted(
candles,
key=lambda candle: candle.open_time,
)
)
def _candle_timestamp_ms_to_utc_datetime(value: int) -> datetime:
try:
return datetime.fromtimestamp(
value / 1000,
tz=timezone.utc,
)
except (OverflowError, OSError, ValueError) as exc:
raise CandleMappingError(
"Поле openTime свечи невозможно преобразовать "
"в UTC datetime."
) from exc
def _required_candle_decimal(
value: DzengiRawNumeric,
*,
field_name: str,
) -> Decimal:
try:
result = Decimal(str(value))
except (InvalidOperation, ValueError) as exc:
raise CandleMappingError(
f"Поле {field_name} свечи невозможно "
"преобразовать в Decimal."
) from exc
if not result.is_finite():
raise CandleMappingError(
f"Поле {field_name} свечи должно быть конечным числом."
)
return result

View File

@@ -131,3 +131,20 @@ class DzengiWebSocketQuoteResponse:
bid_price: DzengiRawNumeric
ask_price: DzengiRawNumeric
timestamp: int | None
# Транспортное представление одной свечи Dzengi GET /api/v1/klines.
@dataclass(frozen=True, slots=True)
class DzengiKline:
open_time: int
open_price: str | int | float
high_price: str | int | float
low_price: str | int | float
close_price: str | int | float
volume: str | int | float
# Транспортное представление ответа Dzengi GET /api/v1/klines.
@dataclass(frozen=True, slots=True)
class DzengiKlinesResponse:
items: tuple[DzengiKline, ...]

View File

@@ -11,6 +11,8 @@ from src.market_data.acquisition.adapters.dzengi.models import (
DzengiInstrumentFilter,
DzengiJsonNumber,
DzengiJsonScalar,
DzengiKline,
DzengiKlinesResponse,
DzengiLotSizeFilter,
DzengiMinNotionalFilter,
DzengiRateLimit,
@@ -20,10 +22,12 @@ from src.market_data.acquisition.adapters.dzengi.models import (
DzengiWebSocketQuoteResponse,
)
from src.market_data.acquisition.exceptions import (
CandleParseError,
InstrumentReferenceParseError,
QuoteParseError,
)
from src.market_data.acquisition.validation.schema import (
ValidatedCandlesDocument,
ValidatedExchangeInfoDocument,
ValidatedQuoteDocument,
ValidatedWebSocketQuoteDocument,
@@ -714,3 +718,153 @@ def _websocket_optional_timestamp(
return None
return _quote_required_int(value, path=path)
# Преобразовать структурно проверенный ответ klines в raw-модели Dzengi.
def parse_candles(
document: ValidatedCandlesDocument,
) -> DzengiKlinesResponse:
"""
Преобразовать элементы проверенного документа klines в transport-модели.
Функция не выполняет предметную проверку числовых значений,
OHLC-инвариантов или mapping во внутреннюю модель Candle.
"""
items = tuple(
_parse_candle_item(
item,
path=f"$.candles[{index}]",
)
for index, item in enumerate(document.items)
)
return DzengiKlinesResponse(items=items)
def _parse_candle_item(
item: object,
*,
path: str,
) -> DzengiKline:
if isinstance(item, Mapping):
return _parse_candle_mapping(item, path=path)
if isinstance(item, tuple):
return _parse_candle_sequence(item, path=path)
raise CandleParseError(
f"{path} должен быть проверенным JSON-объектом или массивом."
)
def _parse_candle_mapping(
item: Mapping[str, object],
*,
path: str,
) -> DzengiKline:
open_time_key = _first_present_key(
item,
keys=("openTime", "open_time", "time", "timestamp"),
)
if open_time_key is None:
raise CandleParseError(
f"{path} не содержит openTime, open_time, time или timestamp."
)
return DzengiKline(
open_time=_candle_required_int(
item.get(open_time_key),
path=f"{path}.{open_time_key}",
),
open_price=_candle_required_raw_numeric(
_required_candle_field(item, key="open", path=path),
path=f"{path}.open",
),
high_price=_candle_required_raw_numeric(
_required_candle_field(item, key="high", path=path),
path=f"{path}.high",
),
low_price=_candle_required_raw_numeric(
_required_candle_field(item, key="low", path=path),
path=f"{path}.low",
),
close_price=_candle_required_raw_numeric(
_required_candle_field(item, key="close", path=path),
path=f"{path}.close",
),
volume=_candle_required_raw_numeric(
_required_candle_field(item, key="volume", path=path),
path=f"{path}.volume",
),
)
def _parse_candle_sequence(
item: tuple[object, ...],
*,
path: str,
) -> DzengiKline:
if len(item) < 6:
raise CandleParseError(
f"{path} должен содержать минимум 6 элементов."
)
return DzengiKline(
open_time=_candle_required_int(item[0], path=f"{path}[0]"),
open_price=_candle_required_raw_numeric(item[1], path=f"{path}[1]"),
high_price=_candle_required_raw_numeric(item[2], path=f"{path}[2]"),
low_price=_candle_required_raw_numeric(item[3], path=f"{path}[3]"),
close_price=_candle_required_raw_numeric(item[4], path=f"{path}[4]"),
volume=_candle_required_raw_numeric(item[5], path=f"{path}[5]"),
)
def _first_present_key(
mapping: Mapping[str, object],
*,
keys: tuple[str, ...],
) -> str | None:
for key in keys:
if key in mapping:
return key
return None
def _required_candle_field(
mapping: Mapping[str, object],
*,
key: str,
path: str,
) -> object:
if key not in mapping:
raise CandleParseError(
f"{path}.{key} отсутствует в элементе klines."
)
return mapping.get(key)
def _candle_required_int(
value: object,
*,
path: str,
) -> int:
if isinstance(value, bool) or not isinstance(value, int):
raise CandleParseError(
f"{path} должен быть целым числом, "
f"получен {type(value).__name__}."
)
return value
def _candle_required_raw_numeric(
value: object,
*,
path: str,
) -> DzengiRawNumeric:
if isinstance(value, bool) or not isinstance(value, (str, int, float)):
raise CandleParseError(
f"{path} должен быть строкой или числом, "
f"получен {type(value).__name__}."
)
return value

View File

@@ -6,6 +6,7 @@ from typing import Protocol
from src.integrations.exchange.rest_client import ExchangeRestClient
from src.market_data.acquisition.exceptions import (
CandleTransportError,
InstrumentReferenceTransportError,
QuoteTransportError,
)
@@ -13,6 +14,7 @@ from src.market_data.acquisition.exceptions import (
_EXCHANGE_INFO_PATH = "/api/v1/exchangeInfo"
_TICKER_24HR_PATH = "/api/v1/ticker/24hr"
_KLINES_PATH = "/api/v1/klines"
# Минимальный транспортный контракт, необходимый Dzengi REST adapter.
@@ -104,3 +106,52 @@ class DzengiQuoteDocumentSource:
"Не удалось получить текущую котировку "
f"от Dzengi для символа '{symbol}': {exc}"
) from exc
class DzengiCandlesDocumentSource:
"""Источник сырого документа свечей через Dzengi REST API."""
def __init__(
self,
client: _PayloadRestClient | None = None,
) -> None:
self._client = client
def fetch_candles_document(
self,
symbol: str,
*,
interval: str,
limit: int,
price_type: str,
) -> object:
"""
Получить декодированный ответ Dzengi klines без его обработки.
Метод не выполняет нормализацию параметров запроса, schema validation,
parsing, value validation, mapping, retry, сортировку или кэширование.
"""
try:
client: _PayloadRestClient = (
self._client
if self._client is not None
else ExchangeRestClient()
)
return client.get_payload(
_KLINES_PATH,
params={
"symbol": symbol,
"interval": interval,
"limit": str(limit),
"priceType": price_type,
},
)
except Exception as exc:
raise CandleTransportError(
"Не удалось получить свечи "
f"от Dzengi для символа '{symbol}', "
f"интервала '{interval}' и типа цены '{price_type}': {exc}"
) from exc

View File

@@ -37,6 +37,7 @@ class InstrumentReferenceMappingError(MarketDataAcquisitionError):
class InstrumentFeedRegistryError(MarketDataAcquisitionError):
pass
# Ошибка получения Quotes Feed от внешнего источника.
class QuoteTransportError(MarketDataAcquisitionError):
pass
@@ -65,3 +66,28 @@ class QuoteMappingError(MarketDataAcquisitionError):
# Ошибка регистрации или получения Quotes Feed.
class QuoteFeedRegistryError(MarketDataAcquisitionError):
pass
# Ошибка получения Candles Feed от внешнего источника.
class CandleTransportError(MarketDataAcquisitionError):
pass
# Ошибка структуры документа Candles Feed.
class CandleSchemaError(MarketDataAcquisitionError):
pass
# Ошибка преобразования проверенного документа в raw-модели свечей.
class CandleParseError(MarketDataAcquisitionError):
pass
# Ошибка допустимости значений Candles Feed.
class CandleValueError(MarketDataAcquisitionError):
pass
# Ошибка преобразования raw-модели источника во внутреннюю модель Candle.
class CandleMappingError(MarketDataAcquisitionError):
pass

View File

@@ -0,0 +1,53 @@
# app/src/market_data/acquisition/feeds/candles_feed.py
from __future__ import annotations
from src.market_data.acquisition.models.candle import Candle
from src.market_data.acquisition.protocol import (
CandlesDocumentHandler,
CandlesDocumentSource,
)
class CandlesFeed:
"""
Read-only Feed канонических рыночных свечей.
Feed координирует получение и обработку документа, но не выполняет
transport parsing, value validation, mapping или кэширование.
"""
def __init__(
self,
*,
source: CandlesDocumentSource,
handler: CandlesDocumentHandler,
) -> None:
self._source = source
self._handler = handler
def load_candles(
self,
symbol: str,
*,
interval: str,
limit: int,
price_type: str,
) -> tuple[Candle, ...]:
"""
Получить канонический immutable-набор свечей инструмента.
"""
document = self._source.fetch_candles_document(
symbol,
interval=interval,
limit=limit,
price_type=price_type,
)
return self._handler.handle_candles_document(
document,
symbol=symbol,
interval=interval,
source=f"rest_klines:{price_type}",
)

View File

@@ -0,0 +1,54 @@
# app/src/market_data/acquisition/handlers/candles_handler.py
from __future__ import annotations
from src.market_data.acquisition.adapters.dzengi.mapper import (
map_dzengi_klines_to_candles,
)
from src.market_data.acquisition.adapters.dzengi.parser import (
parse_candles,
)
from src.market_data.acquisition.models.candle import Candle
from src.market_data.acquisition.validation.schema import (
validate_candles_schema,
)
from src.market_data.acquisition.validation.values import (
validate_candles_values,
)
# Обработчик документа Candles Feed формата Dzengi /api/v1/klines.
class DzengiCandlesDocumentHandler:
def handle_candles_document(
self,
document: object,
*,
symbol: str,
interval: str,
source: str,
) -> tuple[Candle, ...]:
"""
Преобразовать сырой документ Dzengi в канонические модели Candle.
Порядок обработки:
schema validation
parsing
value validation
mapping
"""
validated_document = validate_candles_schema(document)
response = parse_candles(validated_document)
validate_candles_values(response)
return map_dzengi_klines_to_candles(
response,
symbol=symbol,
interval=interval,
source=source,
)

View File

@@ -0,0 +1,24 @@
# app/src/market_data/acquisition/models/candle.py
from __future__ import annotations
from dataclasses import dataclass
from datetime import datetime
from decimal import Decimal
# Каноническая неизменяемая модель одной рыночной свечи OHLCV.
@dataclass(frozen=True, slots=True)
class Candle:
symbol: str
interval: str
open_time: datetime
open_price: Decimal
high_price: Decimal
low_price: Decimal
close_price: Decimal
volume: Decimal
source: str

View File

@@ -4,6 +4,7 @@ from __future__ import annotations
from typing import Protocol, runtime_checkable
from src.market_data.acquisition.models.candle import Candle
from src.market_data.acquisition.models.instrument import Instrument
from src.market_data.acquisition.models.quote import Quote
@@ -83,4 +84,58 @@ class QuoteFeedProtocol(Protocol):
"""
Получить внутреннюю модель текущей котировки инструмента.
"""
...
# Источник сырого документа Candles Feed.
@runtime_checkable
class CandlesDocumentSource(Protocol):
def fetch_candles_document(
self,
symbol: str,
*,
interval: str,
limit: int,
price_type: str,
) -> object:
"""
Получить декодированный транспортный документ рыночных свечей.
Источник не выполняет schema validation, parsing, value validation
или mapping во внутреннюю модель Candle.
"""
...
# Обработчик сырого документа Candles Feed.
@runtime_checkable
class CandlesDocumentHandler(Protocol):
def handle_candles_document(
self,
document: object,
*,
symbol: str,
interval: str,
source: str,
) -> tuple[Candle, ...]:
"""
Преобразовать сырой документ в проверенный immutable-набор Candle.
"""
...
# Источник готового набора свечей.
@runtime_checkable
class CandlesFeedProtocol(Protocol):
def load_candles(
self,
symbol: str,
*,
interval: str,
limit: int,
price_type: str,
) -> tuple[Candle, ...]:
"""
Получить immutable-набор внутренних моделей Candle.
"""
...

View File

@@ -7,6 +7,7 @@ from types import MappingProxyType
from typing import Mapping
from src.market_data.acquisition.exceptions import (
CandleSchemaError,
InstrumentReferenceSchemaError,
QuoteSchemaError,
)
@@ -199,6 +200,7 @@ def _require_list(
return value
# Структурно проверенное представление ответа ticker/24hr.
@dataclass(frozen=True, slots=True)
class ValidatedQuoteDocument:
@@ -398,3 +400,195 @@ def _validate_websocket_quote_payload(
raise QuoteSchemaError(
"WebSocket quote не содержит bid/ask либо bids/asks."
)
# Структурно проверенное представление ответа Dzengi /api/v1/klines.
@dataclass(frozen=True, slots=True)
class ValidatedCandlesDocument:
items: tuple[object, ...]
def validate_candles_schema(
document: object,
) -> ValidatedCandlesDocument:
"""
Проверить структуру ответа Dzengi klines без проверки значений.
Поддерживаются:
1. Корневой JSON-массив:
[
[...],
{...}
]
2. Корневой JSON-объект с массивом в одном из полей:
{
"klines": [...]
}
{
"candles": [...]
}
{
"data": [...]
}
{
"result": [...]
}
3. Wrapped-формат:
{
"payload": [...]
}
{
"payload": {
"klines": [...]
}
}
Внутри payload также поддерживаются поля candles и data.
"""
raw_items = _extract_candle_items(document)
validated_items: list[object] = []
for index, item in enumerate(raw_items):
path = f"$.candles[{index}]"
validated_items.append(
_validate_and_freeze_candle_item(
item,
path=path,
)
)
return ValidatedCandlesDocument(
items=tuple(validated_items),
)
def _extract_candle_items(
document: object,
) -> list[object]:
if isinstance(document, list):
return document
root = _require_candle_mapping(
document,
path="$",
)
direct_items = _find_candle_list(
root,
keys=("klines", "candles", "data", "result"),
path="$",
)
if direct_items is not None:
return direct_items
if "payload" not in root:
raise CandleSchemaError(
"$ не содержит поддерживаемый массив klines."
)
payload = root.get("payload")
if isinstance(payload, list):
return payload
payload_mapping = _require_candle_mapping(
payload,
path="$.payload",
)
wrapped_items = _find_candle_list(
payload_mapping,
keys=("klines", "candles", "data"),
path="$.payload",
)
if wrapped_items is not None:
return wrapped_items
raise CandleSchemaError(
"$.payload не содержит поддерживаемый массив klines."
)
def _find_candle_list(
mapping: Mapping[str, object],
*,
keys: tuple[str, ...],
path: str,
) -> list[object] | None:
for key in keys:
if key not in mapping:
continue
value = mapping.get(key)
if not isinstance(value, list):
raise CandleSchemaError(
f"{path}.{key} должен быть JSON-массивом, "
f"получен {type(value).__name__}."
)
return value
return None
def _validate_and_freeze_candle_item(
item: object,
*,
path: str,
) -> object:
if isinstance(item, dict):
mapping = _require_candle_mapping(
item,
path=path,
)
return MappingProxyType(dict(mapping))
if isinstance(item, list):
if len(item) < 6:
raise CandleSchemaError(
f"{path} должен содержать минимум 6 элементов, "
f"получено {len(item)}."
)
return tuple(item)
raise CandleSchemaError(
f"{path} должен быть JSON-объектом или JSON-массивом, "
f"получен {type(item).__name__}."
)
def _require_candle_mapping(
value: object,
*,
path: str,
) -> Mapping[str, object]:
if not isinstance(value, dict):
raise CandleSchemaError(
f"{path} должен быть JSON-объектом, "
f"получен {type(value).__name__}."
)
for key in value:
if not isinstance(key, str):
raise CandleSchemaError(
f"{path} содержит нестроковый ключ "
f"типа {type(key).__name__}."
)
return value

View File

@@ -13,10 +13,12 @@ from src.market_data.acquisition.adapters.dzengi.models import (
DzengiRateLimit,
DzengiRawNumeric,
DzengiUnknownFilter,
DzengiKlinesResponse,
DzengiTicker24hrResponse,
DzengiWebSocketQuoteResponse,
)
from src.market_data.acquisition.exceptions import (
CandleValueError,
InstrumentReferenceValueError,
QuoteValueError,
)
@@ -454,6 +456,7 @@ def _to_finite_decimal(
return decimal_value
def validate_quote_values(
response: DzengiTicker24hrResponse,
) -> None:
@@ -549,3 +552,128 @@ def validate_dzengi_websocket_quote_values(
raise QuoteValueError(
"$.payload.timestamp должно быть больше нуля."
)
def validate_candles_values(
response: DzengiKlinesResponse,
) -> None:
"""
Проверить допустимость значений raw-моделей Dzengi klines.
Функция не изменяет transport-модели, не выполняет mapping в Candle
и не проверяет последовательность временных меток между свечами.
"""
for index, candle in enumerate(response.items):
path = f"$.candles[{index}]"
if isinstance(candle.open_time, bool) or candle.open_time <= 0:
raise CandleValueError(
f"{path}.openTime должно быть целым числом больше нуля."
)
open_price = _candle_positive_decimal(
candle.open_price,
path=f"{path}.open",
)
high_price = _candle_positive_decimal(
candle.high_price,
path=f"{path}.high",
)
low_price = _candle_positive_decimal(
candle.low_price,
path=f"{path}.low",
)
close_price = _candle_positive_decimal(
candle.close_price,
path=f"{path}.close",
)
_candle_non_negative_decimal(
candle.volume,
path=f"{path}.volume",
)
if high_price < low_price:
raise CandleValueError(
f"{path}.high не должно быть меньше {path}.low."
)
if high_price < open_price:
raise CandleValueError(
f"{path}.high не должно быть меньше {path}.open."
)
if high_price < close_price:
raise CandleValueError(
f"{path}.high не должно быть меньше {path}.close."
)
if low_price > open_price:
raise CandleValueError(
f"{path}.low не должно превышать {path}.open."
)
if low_price > close_price:
raise CandleValueError(
f"{path}.low не должно превышать {path}.close."
)
def _candle_positive_decimal(
value: object,
*,
path: str,
) -> Decimal:
decimal_value = _candle_decimal(
value,
path=path,
)
if decimal_value <= 0:
raise CandleValueError(
f"{path} должно быть больше нуля."
)
return decimal_value
def _candle_non_negative_decimal(
value: object,
*,
path: str,
) -> Decimal:
decimal_value = _candle_decimal(
value,
path=path,
)
if decimal_value < 0:
raise CandleValueError(
f"{path} должно быть больше или равно нулю."
)
return decimal_value
def _candle_decimal(
value: object,
*,
path: str,
) -> Decimal:
if isinstance(value, bool) or not isinstance(value, (str, int, float)):
raise CandleValueError(
f"{path} должно быть числом или числовой строкой."
)
try:
decimal_value = Decimal(str(value))
except (InvalidOperation, ValueError) as exc:
raise CandleValueError(
f"{path} должно быть корректным числом."
) from exc
if not decimal_value.is_finite():
raise CandleValueError(
f"{path} должно быть конечным числом."
)
return decimal_value