From 77e87c0504bb843224450e023b1fe9117a0dec0d Mon Sep 17 00:00:00 2001 From: Sergey Date: Wed, 15 Jul 2026 16:35:54 +0300 Subject: [PATCH] build 047: switch ExchangeService klines to canonical candles feed --- app/src/integrations/exchange/service.py | 199 ++---- .../exchange/test_service_klines.py | 551 +++++++++++++++++ docs/migrations/build_047.md | 573 ++++++++++++++++++ 3 files changed, 1184 insertions(+), 139 deletions(-) create mode 100644 app/tests/unit/integrations/exchange/test_service_klines.py create mode 100644 docs/migrations/build_047.md diff --git a/app/src/integrations/exchange/service.py b/app/src/integrations/exchange/service.py index 32f1352..8f56296 100644 --- a/app/src/integrations/exchange/service.py +++ b/app/src/integrations/exchange/service.py @@ -41,24 +41,32 @@ from src.integrations.exchange.status import ( classify_exchange_error, ) from src.market_data.acquisition.adapters.dzengi.rest import ( + DzengiCandlesDocumentSource, DzengiInstrumentDocumentSource, DzengiQuoteDocumentSource, ) +from src.market_data.acquisition.feeds.candles_feed import CandlesFeed from src.market_data.acquisition.feeds.instrument_feed import InstrumentFeed from src.market_data.acquisition.feeds.quotes_feed import QuotesFeed +from src.market_data.acquisition.handlers.candles_handler import ( + DzengiCandlesDocumentHandler, +) from src.market_data.acquisition.handlers.instrument_handler import ( DzengiInstrumentDocumentHandler, ) from src.market_data.acquisition.handlers.quotes_handler import ( DzengiQuoteDocumentHandler, ) +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 from src.market_data.acquisition.registry import ( + CandlesFeedRegistry, InstrumentFeedRegistry, QuoteFeedRegistry, ) from src.market_data.acquisition.service import ( + CandlesAcquisitionService, InstrumentAcquisitionService, QuoteAcquisitionService, ) @@ -79,6 +87,7 @@ from src.trading.journal.service import JournalService _INSTRUMENT_REFERENCE_SOURCE_NAME = "dzengi" _QUOTE_SOURCE_NAME = "dzengi" +_CANDLES_SOURCE_NAME = "dzengi" class ExchangeService: @@ -491,7 +500,7 @@ class ExchangeService: def _runtime_key(self, runtime_key: str | None) -> str: return (runtime_key or self._default_runtime_key).strip().lower() - # Получить свечи инструмента через REST API. + # Получить свечи инструмента через канонический Candles Feed. def get_klines( self, symbol: str | None = None, @@ -523,17 +532,12 @@ class ExchangeService: if not validation.is_valid: raise ExchangeError(validation.message) - client = ExchangeRestClient() - try: - payload = client.get_payload( - "/api/v1/klines", - params={ - "symbol": validation.normalized_symbol, - "interval": interval, - "limit": str(limit), - "priceType": normalized_price_type, - }, + canonical_candles = self._load_candles_via_acquisition( + symbol=validation.normalized_symbol, + interval=interval, + limit=limit, + price_type=normalized_price_type, ) except Exception as exc: self._log_exchange_error( @@ -548,149 +552,66 @@ class ExchangeService: ) raise ExchangeError(str(exc)) from exc - candles = self._parse_klines_payload( - payload=payload, - symbol=validation.normalized_symbol, - interval=interval, - source=f"rest_klines:{normalized_price_type}", - ) + candles = [ + self._kline_from_candle(candle) + for candle in canonical_candles[-limit:] + ] return KlineBatch( symbol=validation.normalized_symbol, interval=interval, - candles=candles[-limit:], + candles=candles, source=f"rest_klines:{normalized_price_type}", ) - # Преобразовать сырой payload свечей в список Kline. - def _parse_klines_payload( + # Собрать Candles acquisition pipeline и вернуть канонические модели Candle. + def _load_candles_via_acquisition( self, *, - payload: object, symbol: str, interval: str, - source: str, - ) -> list[Kline]: - raw_items = self._extract_klines_items(payload) + limit: int, + price_type: str, + ) -> tuple[Candle, ...]: + source = DzengiCandlesDocumentSource() + handler = DzengiCandlesDocumentHandler() - candles: list[Kline] = [] + feed = CandlesFeed( + source=source, + handler=handler, + ) - for item in raw_items: - candle = self._parse_kline_item( - item=item, - symbol=symbol, - interval=interval, - source=source, - ) + registry = CandlesFeedRegistry() + registry.register( + _CANDLES_SOURCE_NAME, + feed, + ) - if candle is not None: - candles.append(candle) + acquisition_service = CandlesAcquisitionService( + registry=registry, + ) - candles.sort(key=lambda item: item.open_time) + return acquisition_service.load_candles( + _CANDLES_SOURCE_NAME, + symbol, + interval=interval, + limit=limit, + price_type=price_type, + ) - return candles - - # Извлечь массив свечей из разных возможных форматов ответа API. - def _extract_klines_items(self, payload: object) -> list[object]: - if isinstance(payload, list): - return payload - - if not isinstance(payload, dict): - return [] - - for key in ("klines", "candles", "data", "result"): - value = payload.get(key) - if isinstance(value, list): - return value - - inner = payload.get("payload") - - if isinstance(inner, list): - return inner - - if isinstance(inner, dict): - for key in ("klines", "candles", "data"): - value = inner.get(key) - if isinstance(value, list): - return value - - return [] - - # Преобразовать одну свечу из dict/list формата в Kline. - def _parse_kline_item( - self, - *, - item: object, - symbol: str, - interval: str, - source: str, - ) -> Kline | None: - if isinstance(item, dict): - open_time = ( - item.get("openTime") - or item.get("open_time") - or item.get("time") - or item.get("timestamp") - ) - - open_time_value = safe_float(open_time) - open_price = safe_float(item.get("open")) - high_price = safe_float(item.get("high")) - low_price = safe_float(item.get("low")) - close_price = safe_float(item.get("close")) - volume = safe_float(item.get("volume")) or 0.0 - - if ( - open_time_value is None - or open_price is None - or high_price is None - or low_price is None - or close_price is None - ): - return None - - return Kline( - symbol=symbol, - interval=interval, - open_time=int(open_time_value), - open_price=open_price, - high_price=high_price, - low_price=low_price, - close_price=close_price, - volume=volume, - source=source, - ) - - if isinstance(item, list) and len(item) >= 6: - open_time_value = safe_float(item[0]) - open_price = safe_float(item[1]) - high_price = safe_float(item[2]) - low_price = safe_float(item[3]) - close_price = safe_float(item[4]) - volume = safe_float(item[5]) or 0.0 - - if ( - open_time_value is None - or open_price is None - or high_price is None - or low_price is None - or close_price is None - ): - return None - - return Kline( - symbol=symbol, - interval=interval, - open_time=int(open_time_value), - open_price=open_price, - high_price=high_price, - low_price=low_price, - close_price=close_price, - volume=volume, - source=source, - ) - - return None + # Временно преобразовать canonical Candle в legacy Kline. + def _kline_from_candle(self, candle: Candle) -> Kline: + return Kline( + symbol=candle.symbol, + interval=candle.interval, + open_time=int(candle.open_time.timestamp() * 1000), + open_price=float(candle.open_price), + high_price=float(candle.high_price), + low_price=float(candle.low_price), + close_price=float(candle.close_price), + volume=float(candle.volume), + source=candle.source, + ) # Проверить публичную доступность биржи. def get_health(self) -> ExchangeHealth: @@ -1244,4 +1165,4 @@ class ExchangeService: except Exception: pass - return None \ No newline at end of file + return None diff --git a/app/tests/unit/integrations/exchange/test_service_klines.py b/app/tests/unit/integrations/exchange/test_service_klines.py new file mode 100644 index 0000000..a0e9f39 --- /dev/null +++ b/app/tests/unit/integrations/exchange/test_service_klines.py @@ -0,0 +1,551 @@ +# app/tests/unit/integrations/exchange/test_service_klines.py + +from __future__ import annotations + +from datetime import datetime, timezone +from decimal import Decimal +from types import SimpleNamespace +from typing import Any, cast + +import pytest + +import src.integrations.exchange.service as service_module +from src.integrations.exchange.exceptions import ExchangeError +from src.integrations.exchange.models import SymbolValidationResult +from src.integrations.exchange.service import ExchangeService +from src.market_data.acquisition.models.candle import Candle + + +def _candle( + *, + symbol: str = "BTC/USD_LEVERAGE", + interval: str = "1m", + open_time_ms: int = 1_750_000_000_000, + open_price: str = "100.10", + high_price: str = "110.20", + low_price: str = "90.30", + close_price: str = "105.40", + volume: str = "12.50", + source: str = "rest_klines:bid", +) -> Candle: + return Candle( + symbol=symbol, + interval=interval, + open_time=datetime.fromtimestamp( + open_time_ms / 1000, + tz=timezone.utc, + ), + open_price=Decimal(open_price), + high_price=Decimal(high_price), + low_price=Decimal(low_price), + close_price=Decimal(close_price), + volume=Decimal(volume), + source=source, + ) + + + +def _set_test_attribute( + target: object, + name: str, + value: object, +) -> None: + setattr(cast(Any, target), name, value) + + +def _service( + *, + default_symbol: str = "BTC/USD_LEVERAGE", + exchange_enabled: bool = True, +) -> ExchangeService: + service = ExchangeService.__new__(ExchangeService) + _set_test_attribute( + service, + "settings", + SimpleNamespace( + default_symbol=default_symbol, + exchange_enabled=exchange_enabled, + ), + ) + return service + + +def _valid_symbol( + symbol: str = "BTC/USD_LEVERAGE", +) -> SymbolValidationResult: + return SymbolValidationResult( + requested_symbol=symbol, + normalized_symbol=symbol, + is_valid=True, + message="OK", + symbol_info=None, + ) + + +def _invalid_symbol( + symbol: str = "UNKNOWN", +) -> SymbolValidationResult: + return SymbolValidationResult( + requested_symbol=symbol, + normalized_symbol=symbol, + is_valid=False, + message="Invalid symbol.", + symbol_info=None, + ) + + +def test_get_klines_uses_default_symbol() -> None: + service = _service(default_symbol="ETH/USD_LEVERAGE") + requested_symbols: list[str] = [] + acquisition_calls: list[dict[str, object]] = [] + + _set_test_attribute( + service, + "validate_symbol", + lambda symbol: ( + requested_symbols.append(symbol) or _valid_symbol(symbol) + ), + ) + _set_test_attribute( + service, + "_load_candles_via_acquisition", + lambda **kwargs: acquisition_calls.append(kwargs) or (), + ) + + result = service.get_klines() + + assert requested_symbols == ["ETH/USD_LEVERAGE"] + assert acquisition_calls[0]["symbol"] == "ETH/USD_LEVERAGE" + assert result.symbol == "ETH/USD_LEVERAGE" + + +@pytest.mark.parametrize("limit", [0, -1, -100]) +def test_get_klines_normalizes_non_positive_limit(limit: int) -> None: + service = _service() + captured: list[dict[str, object]] = [] + + _set_test_attribute( + service, + "validate_symbol", + lambda symbol: _valid_symbol(symbol), + ) + _set_test_attribute( + service, + "_load_candles_via_acquisition", + lambda **kwargs: captured.append(kwargs) or (), + ) + + service.get_klines(limit=limit) + + assert captured[0]["limit"] == 200 + + +def test_get_klines_caps_limit_at_200() -> None: + service = _service() + captured: list[dict[str, object]] = [] + + _set_test_attribute( + service, + "validate_symbol", + lambda symbol: _valid_symbol(symbol), + ) + _set_test_attribute( + service, + "_load_candles_via_acquisition", + lambda **kwargs: captured.append(kwargs) or (), + ) + + service.get_klines(limit=500) + + assert captured[0]["limit"] == 200 + + +@pytest.mark.parametrize("interval", ["1m", "5m", "15m", "1h"]) +def test_get_klines_accepts_supported_intervals(interval: str) -> None: + service = _service() + captured: list[dict[str, object]] = [] + + _set_test_attribute( + service, + "validate_symbol", + lambda symbol: _valid_symbol(symbol), + ) + _set_test_attribute( + service, + "_load_candles_via_acquisition", + lambda **kwargs: captured.append(kwargs) or (), + ) + + result = service.get_klines(interval=interval) + + assert captured[0]["interval"] == interval + assert result.interval == interval + + +def test_get_klines_rejects_unsupported_interval() -> None: + service = _service() + + with pytest.raises( + ExchangeError, + match="Unsupported kline interval", + ): + service.get_klines(interval="4h") + + +@pytest.mark.parametrize( + ("price_type", "expected"), + [ + ("bid", "bid"), + ("ask", "ask"), + (" BID ", "bid"), + (" AsK ", "ask"), + ("unknown", "bid"), + ("", "bid"), + ], +) +def test_get_klines_normalizes_price_type( + price_type: str, + expected: str, +) -> None: + service = _service() + captured: list[dict[str, object]] = [] + + _set_test_attribute( + service, + "validate_symbol", + lambda symbol: _valid_symbol(symbol), + ) + _set_test_attribute( + service, + "_load_candles_via_acquisition", + lambda **kwargs: captured.append(kwargs) or (), + ) + + result = service.get_klines(price_type=price_type) + + assert captured[0]["price_type"] == expected + assert result.source == f"rest_klines:{expected}" + + +def test_get_klines_rejects_mock_mode() -> None: + service = _service(exchange_enabled=False) + + with pytest.raises( + ExchangeError, + match="Klines are not available in mock exchange mode", + ): + service.get_klines() + + +def test_get_klines_rejects_invalid_symbol() -> None: + service = _service() + _set_test_attribute( + service, + "validate_symbol", + lambda symbol: _invalid_symbol(symbol), + ) + + with pytest.raises(ExchangeError, match="Invalid symbol"): + service.get_klines("UNKNOWN") + + +def test_get_klines_calls_acquisition_with_exact_arguments() -> None: + service = _service() + captured: list[dict[str, object]] = [] + + _set_test_attribute( + service, + "validate_symbol", + lambda symbol: _valid_symbol("BTC/USD_LEVERAGE"), + ) + _set_test_attribute( + service, + "_load_candles_via_acquisition", + lambda **kwargs: captured.append(kwargs) or (), + ) + + service.get_klines( + " btc/usd ", + interval="5m", + limit=123, + price_type="ask", + ) + + assert captured == [ + { + "symbol": "BTC/USD_LEVERAGE", + "interval": "5m", + "limit": 123, + "price_type": "ask", + } + ] + + +def test_get_klines_maps_canonical_candle_to_legacy_kline() -> None: + service = _service() + candle = _candle() + + _set_test_attribute( + service, + "validate_symbol", + lambda symbol: _valid_symbol(symbol), + ) + _set_test_attribute( + service, + "_load_candles_via_acquisition", + lambda **kwargs: (candle,), + ) + + result = service.get_klines() + + assert len(result.candles) == 1 + kline = result.candles[0] + assert kline.symbol == candle.symbol + assert kline.interval == candle.interval + assert kline.open_time == 1_750_000_000_000 + assert kline.open_price == 100.10 + assert kline.high_price == 110.20 + assert kline.low_price == 90.30 + assert kline.close_price == 105.40 + assert kline.volume == 12.50 + assert kline.source == candle.source + + +def test_get_klines_preserves_candle_order() -> None: + service = _service() + first = _candle(open_time_ms=1000) + second = _candle(open_time_ms=2000) + + _set_test_attribute( + service, + "validate_symbol", + lambda symbol: _valid_symbol(symbol), + ) + _set_test_attribute( + service, + "_load_candles_via_acquisition", + lambda **kwargs: (first, second), + ) + + result = service.get_klines() + + assert [item.open_time for item in result.candles] == [1000, 2000] + + +def test_get_klines_trims_result_to_limit() -> None: + service = _service() + candles = tuple( + _candle(open_time_ms=index * 1000) + for index in range(1, 6) + ) + + _set_test_attribute( + service, + "validate_symbol", + lambda symbol: _valid_symbol(symbol), + ) + _set_test_attribute( + service, + "_load_candles_via_acquisition", + lambda **kwargs: candles, + ) + + result = service.get_klines(limit=2) + + assert [item.open_time for item in result.candles] == [4000, 5000] + + +def test_get_klines_returns_empty_batch() -> None: + service = _service() + + _set_test_attribute( + service, + "validate_symbol", + lambda symbol: _valid_symbol(symbol), + ) + _set_test_attribute( + service, + "_load_candles_via_acquisition", + lambda **kwargs: (), + ) + + result = service.get_klines() + + assert result.candles == [] + assert result.source == "rest_klines:bid" + + +def test_get_klines_logs_and_wraps_acquisition_error() -> None: + service = _service() + original_error = RuntimeError("candles unavailable") + log_calls: list[dict[str, object]] = [] + + _set_test_attribute( + service, + "validate_symbol", + lambda symbol: _valid_symbol(symbol), + ) + + def raise_error(**kwargs: object) -> tuple[Candle, ...]: + raise original_error + + _set_test_attribute( + service, + "_load_candles_via_acquisition", + raise_error, + ) + _set_test_attribute( + service, + "_log_exchange_error", + lambda **kwargs: log_calls.append(kwargs), + ) + + with pytest.raises( + ExchangeError, + match="candles unavailable", + ) as error_info: + service.get_klines( + interval="15m", + limit=25, + price_type="ask", + ) + + assert error_info.value.__cause__ is original_error + assert log_calls == [ + { + "endpoint": "klines", + "exc": original_error, + "symbol": "BTC/USD_LEVERAGE", + "extra_payload": { + "interval": "15m", + "limit": 25, + "price_type": "ask", + }, + } + ] + + +def test_kline_from_candle_returns_new_legacy_model() -> None: + service = _service() + candle = _candle() + + result = service._kline_from_candle(candle) + + assert result is not candle + assert result.symbol == candle.symbol + assert result.open_time == 1_750_000_000_000 + + +def test_load_candles_via_acquisition_builds_complete_pipeline( + monkeypatch: pytest.MonkeyPatch, +) -> None: + events: list[tuple[str, object]] = [] + expected = (_candle(),) + + class FakeSource: + def __init__(self) -> None: + events.append(("source", self)) + + class FakeHandler: + def __init__(self) -> None: + events.append(("handler", self)) + + class FakeFeed: + def __init__(self, *, source: object, handler: object) -> None: + self.source = source + self.handler = handler + events.append(("feed", self)) + + class FakeRegistry: + def __init__(self) -> None: + self.registered: tuple[str, object] | None = None + events.append(("registry", self)) + + def register(self, source_name: str, feed: object) -> None: + self.registered = (source_name, feed) + events.append(("register", self.registered)) + + class FakeAcquisitionService: + def __init__(self, *, registry: FakeRegistry) -> None: + self.registry = registry + events.append(("service", registry)) + + def load_candles( + self, + source_name: str, + symbol: str, + *, + interval: str, + limit: int, + price_type: str, + ) -> tuple[Candle, ...]: + events.append( + ( + "load", + ( + source_name, + symbol, + interval, + limit, + price_type, + ), + ) + ) + return expected + + monkeypatch.setattr( + service_module, + "DzengiCandlesDocumentSource", + FakeSource, + ) + monkeypatch.setattr( + service_module, + "DzengiCandlesDocumentHandler", + FakeHandler, + ) + monkeypatch.setattr(service_module, "CandlesFeed", FakeFeed) + monkeypatch.setattr( + service_module, + "CandlesFeedRegistry", + FakeRegistry, + ) + monkeypatch.setattr( + service_module, + "CandlesAcquisitionService", + FakeAcquisitionService, + ) + + service = _service() + result = service._load_candles_via_acquisition( + symbol="BTC/USD_LEVERAGE", + interval="1m", + limit=100, + price_type="bid", + ) + + assert result is expected + + source = events[0][1] + handler = events[1][1] + feed = events[2][1] + registry = events[3][1] + + assert isinstance(source, FakeSource) + assert isinstance(handler, FakeHandler) + assert isinstance(feed, FakeFeed) + assert isinstance(registry, FakeRegistry) + + assert feed.source is source + assert feed.handler is handler + assert registry.registered == ("dzengi", feed) + + assert ( + "load", + ( + "dzengi", + "BTC/USD_LEVERAGE", + "1m", + 100, + "bid", + ), + ) in events diff --git a/docs/migrations/build_047.md b/docs/migrations/build_047.md new file mode 100644 index 0000000..23372b6 --- /dev/null +++ b/docs/migrations/build_047.md @@ -0,0 +1,573 @@ +# Build 047 — Переключение `ExchangeService.get_klines()` на канонический Candles Feed + +**Статус:** Completed +**Дата:** 2026-07-15 +**Подсистема:** Market Data Acquisition / OHLCV Feed +**Тип изменения:** Миграция legacy-потребителя на канонический Acquisition pipeline + +--- + +## 1. Цель Build + +Цель Build 047 — переключить существующий публичный метод: + +```text +ExchangeService.get_klines() +``` + +с прямого получения и самостоятельной обработки ответа Dzengi `/api/v1/klines` на канонический pipeline подсистемы: + +```text +src/market_data/acquisition/ +``` + +При этом необходимо сохранить существующий публичный контракт `ExchangeService.get_klines()` и обратную совместимость с текущими потребителями legacy-слоя. + +--- + +## 2. Исходное состояние + +До Build 047 метод: + +```text +ExchangeService.get_klines() +``` + +самостоятельно: + +1. выполнял прямой REST-запрос к `/api/v1/klines`; +2. извлекал элементы ответа; +3. разбирал отдельные свечи; +4. создавал legacy-модели `Kline`; +5. формировал `KlineBatch`. + +В `ExchangeService` существовали собственные legacy-функции обработки свечей: + +```text +_parse_klines_payload() +_extract_klines_items() +_parse_kline_item() +``` + +Таким образом, после появления канонического Candles Feed в системе существовали два независимых пути получения и обработки OHLCV-данных. + +Это создавало дублирование ответственности между: + +```text +src/integrations/exchange/ +``` + +и: + +```text +src/market_data/acquisition/ +``` + +--- + +## 3. Границы Build + +В Build 047 выполнено: + +1. переключение внутренней реализации `ExchangeService.get_klines()` на канонический Candles Feed; +2. сохранение публичного контракта `ExchangeService.get_klines()`; +3. сохранение legacy-моделей `Kline` и `KlineBatch`; +4. добавление внутреннего адаптационного преобразования `Candle -> Kline`; +5. удаление прямого REST-запроса `/api/v1/klines` из `ExchangeService`; +6. удаление legacy-функций: + - `_parse_klines_payload()`; + - `_extract_klines_items()`; + - `_parse_kline_item()`; +7. добавление специализированных unit-тестов для `ExchangeService.get_klines()` и его интеграции с каноническим Acquisition pipeline. + +В Build 047 не выполнялось: + +- изменение публичного контракта `ExchangeService.get_klines()`; +- изменение legacy-потребителей в `trading/market_analysis`; +- удаление моделей `Kline` и `KlineBatch`; +- прямое переключение `trading/market_analysis` на `CandlesAcquisitionService`; +- изменение архитектуры других Market Data feeds; +- изменение runtime-компонентов Acquisition. + +--- + +## 4. Изменённые production-файлы + +### 4.1. `src/integrations/exchange/service.py` + +Метод: + +```text +ExchangeService.get_klines() +``` + +сохранён как публичный compatibility-контракт. + +Внутренний путь получения данных изменён. + +До Build 047: + +```text +ExchangeService.get_klines() + ↓ +прямой REST-запрос /api/v1/klines + ↓ +_parse_klines_payload() + ↓ +_extract_klines_items() + ↓ +_parse_kline_item() + ↓ +Kline + ↓ +KlineBatch +``` + +После Build 047: + +```text +ExchangeService.get_klines() + ↓ +_load_candles_via_acquisition() + ↓ +DzengiCandlesDocumentSource + ↓ +DzengiCandlesDocumentHandler + ↓ +CandlesFeed + ↓ +CandlesFeedRegistry + ↓ +CandlesAcquisitionService + ↓ +Candle + ↓ +_kline_from_candle() + ↓ +Kline + ↓ +KlineBatch +``` + +--- + +## 5. Канонический путь получения свечей + +После Build 047 единственный production-путь прямого обращения к endpoint: + +```text +/api/v1/klines +``` + +расположен в: + +```text +src/market_data/acquisition/adapters/dzengi/rest.py +``` + +Это соответствует архитектурному разделению ответственности: + +```text +Exchange-specific transport + ↓ +Market Data Acquisition + ↓ +Canonical Candle + ↓ +Legacy compatibility boundary + ↓ +Kline / KlineBatch +``` + +`ExchangeService` больше не владеет транспортной логикой получения OHLCV-данных от Dzengi. + +--- + +## 6. Добавленные внутренние методы + +### 6.1. `_load_candles_via_acquisition()` + +Добавлен внутренний compatibility-helper: + +```text +ExchangeService._load_candles_via_acquisition() +``` + +Его ответственность: + +1. создать `DzengiCandlesDocumentSource`; +2. создать `DzengiCandlesDocumentHandler`; +3. создать `CandlesFeed`; +4. зарегистрировать feed в `CandlesFeedRegistry`; +5. создать `CandlesAcquisitionService`; +6. вызвать канонический `load_candles()`; +7. вернуть immutable-набор моделей `Candle`. + +Метод является внутренней границей между legacy `ExchangeService` и канонической подсистемой Market Data Acquisition. + +### 6.2. `_kline_from_candle()` + +Добавлен внутренний адаптер: + +```text +ExchangeService._kline_from_candle() +``` + +Преобразование выполняется по следующему правилу: + +```text +Candle.symbol -> Kline.symbol +Candle.interval -> Kline.interval +Candle.open_time -> Kline.open_time +Candle.open_price -> Kline.open_price +Candle.high_price -> Kline.high_price +Candle.low_price -> Kline.low_price +Candle.close_price -> Kline.close_price +Candle.volume -> Kline.volume +``` + +Это позволяет сохранить существующий legacy-контракт без переноса legacy-моделей внутрь `market_data/acquisition`. + +--- + +## 7. Удалённая legacy-логика + +Из: + +```text +src/integrations/exchange/service.py +``` + +удалены: + +```text +_parse_klines_payload() +_extract_klines_items() +_parse_kline_item() +``` + +Также удалён прямой вызов: + +```text +/api/v1/klines +``` + +из `ExchangeService`. + +После Build 047 parsing, schema validation, value validation и mapping OHLCV-данных выполняются только канонической подсистемой: + +```text +src/market_data/acquisition/ +``` + +--- + +## 8. Сохранённая обратная совместимость + +Публичный контракт: + +```text +ExchangeService.get_klines(...) +``` + +сохранён. + +Существующие production-потребители не изменялись: + +```text +src/trading/market_analysis/service.py +src/trading/market_analysis/htf.py +``` + +Они продолжают работать через: + +```text +ExchangeService.get_klines() +``` + +и получать: + +```text +KlineBatch +``` + +Таким образом, Build 047 меняет внутренний источник данных, но не требует одновременного переписывания legacy-потребителей. + +--- + +## 9. Добавленные тесты + +Добавлен специализированный файл: + +```text +tests/unit/integrations/exchange/test_service_klines.py +``` + +Тесты покрывают: + +- использование `default_symbol`; +- нормализацию `limit`; +- нормализацию `interval`; +- отклонение неподдерживаемого interval; +- передачу `price_type`; +- обработку выключенной биржи; +- обработку неизвестного символа; +- передачу параметров в Acquisition pipeline; +- преобразование `Candle` в `Kline`; +- формирование `KlineBatch`; +- обработку пустого набора свечей; +- обработку исключений канонического Acquisition pipeline; +- построение полного composition pipeline; +- сохранение публичного legacy-контракта. + +Результат targeted-проверки: + +```text +26 passed in 0.17s +``` + +--- + +## 10. Полная регрессионная проверка + +Выполнено: + +```bash +python -m pytest -q +``` + +Результат: + +```text +750 passed in 2.83s +``` + +Регрессий не обнаружено. + +--- + +## 11. Контроль прямого использования `/api/v1/klines` + +Команда: + +```bash +grep -RIn \ + --exclude-dir="__pycache__" \ + --exclude="*.pyc" \ + '"/api/v1/klines"' \ + src +``` + +Результат: + +```text +src/market_data/acquisition/adapters/dzengi/rest.py:17:_KLINES_PATH = "/api/v1/klines" +``` + +Вывод: + +Прямое знание endpoint `/api/v1/klines` находится только в каноническом Dzengi REST-адаптере. + +--- + +## 12. Контроль удаления legacy-парсеров + +Команда: + +```bash +grep -RIn \ + --exclude-dir="__pycache__" \ + --exclude="*.pyc" \ + "_parse_klines_payload\|_extract_klines_items\|_parse_kline_item" \ + src tests +``` + +Результат: + +```text +пусто +``` + +Вывод: + +Legacy-функции обработки kline payload полностью удалены. + +--- + +## 13. Контроль оставшихся потребителей `ExchangeService.get_klines()` + +Команда: + +```bash +grep -RIn \ + --exclude-dir="__pycache__" \ + --exclude="*.pyc" \ + "\.get_klines(" \ + src tests +``` + +Production-потребители: + +```text +src/trading/market_analysis/service.py +src/trading/market_analysis/htf.py +``` + +Дополнительно присутствуют специализированные вызовы в: + +```text +tests/unit/integrations/exchange/test_service_klines.py +``` + +Вывод: + +Публичный compatibility-контракт `ExchangeService.get_klines()` продолжает обслуживать существующие legacy-потребители. + +--- + +## 14. Контроль интеграции с Acquisition pipeline + +Команда: + +```bash +grep -RIn \ + --exclude-dir="__pycache__" \ + --exclude="*.pyc" \ + "CandlesAcquisitionService\|CandlesFeedRegistry\|DzengiCandlesDocumentSource" \ + src/integrations/exchange/service.py +``` + +Результат подтверждает использование: + +```text +DzengiCandlesDocumentSource +CandlesFeedRegistry +CandlesAcquisitionService +``` + +внутри compatibility boundary `ExchangeService`. + +--- + +## 15. Контроль качества diff + +Выполнено: + +```bash +git diff --check +``` + +Результат: + +```text +пусто +``` + +Whitespace-ошибок не обнаружено. + +--- + +## 16. Архитектурный результат + +После Build 047 путь OHLCV-данных имеет следующую структуру: + +```text +Dzengi REST API + ↓ +DzengiCandlesDocumentSource + ↓ +schema validation + ↓ +parser + ↓ +value validation + ↓ +mapper + ↓ +Candle + ↓ +CandlesFeed + ↓ +CandlesFeedRegistry + ↓ +CandlesAcquisitionService + ↓ +ExchangeService compatibility boundary + ↓ +Kline / KlineBatch + ↓ +legacy trading consumers +``` + +Главный архитектурный результат: + +```text +ExchangeService больше не получает и не разбирает +сырой ответ /api/v1/klines самостоятельно. +``` + +Владение получением и канонической обработкой OHLCV-данных теперь принадлежит: + +```text +src/market_data/acquisition/ +``` + +--- + +## 17. Состояние после Build 047 + +На момент завершения Build: + +```text +Targeted tests: 26 passed +Full test suite: 750 passed +git diff --check: clean +``` + +Прямой endpoint: + +```text +/api/v1/klines +``` + +остался только в: + +```text +src/market_data/acquisition/adapters/dzengi/rest.py +``` + +Legacy parsing helpers: + +```text +_parse_klines_payload +_extract_klines_items +_parse_kline_item +``` + +полностью отсутствуют. + +Публичный контракт: + +```text +ExchangeService.get_klines() +``` + +сохранён. + +--- + +## 18. Итог Build + +**Build 047 завершён успешно.** + +Канонический Candles Feed интегрирован в существующий `ExchangeService.get_klines()` без изменения его публичного контракта. + +Дублирующая legacy-логика получения и разбора `/api/v1/klines` удалена. + +Существующие потребители продолжают работать без изменений. + +Полная регрессионная проверка подтверждает отсутствие нарушений существующего поведения: + +```text +750 passed +``` \ No newline at end of file