build 049: expose canonical candles through ExchangeService

This commit is contained in:
2026-07-15 20:54:24 +03:00
parent 3a253d89a9
commit 16ed64f7c6
3 changed files with 1010 additions and 7 deletions

View File

@@ -500,15 +500,15 @@ class ExchangeService:
def _runtime_key(self, runtime_key: str | None) -> str: def _runtime_key(self, runtime_key: str | None) -> str:
return (runtime_key or self._default_runtime_key).strip().lower() return (runtime_key or self._default_runtime_key).strip().lower()
# Получить свечи инструмента через канонический Candles Feed. # Получить канонические свечи инструмента через Candles Feed.
def get_klines( def get_candles(
self, self,
symbol: str | None = None, symbol: str | None = None,
*, *,
interval: str = "1m", interval: str = "1m",
limit: int = 200, limit: int = 200,
price_type: str = "bid", price_type: str = "bid",
) -> KlineBatch: ) -> tuple[Candle, ...]:
symbol_to_use = symbol or self.settings.default_symbol symbol_to_use = symbol or self.settings.default_symbol
if limit <= 0: if limit <= 0:
@@ -526,14 +526,16 @@ class ExchangeService:
normalized_price_type = "bid" normalized_price_type = "bid"
if not self.settings.exchange_enabled: if not self.settings.exchange_enabled:
raise ExchangeError("Klines are not available in mock exchange mode.") raise ExchangeError(
"Candles are not available in mock exchange mode."
)
validation = self.validate_symbol(symbol_to_use) validation = self.validate_symbol(symbol_to_use)
if not validation.is_valid: if not validation.is_valid:
raise ExchangeError(validation.message) raise ExchangeError(validation.message)
try: try:
canonical_candles = self._load_candles_via_acquisition( return self._load_candles_via_acquisition(
symbol=validation.normalized_symbol, symbol=validation.normalized_symbol,
interval=interval, interval=interval,
limit=limit, limit=limit,
@@ -552,13 +554,54 @@ class ExchangeService:
) )
raise ExchangeError(str(exc)) from exc raise ExchangeError(str(exc)) from exc
# Сохранить legacy KlineBatch как compatibility-фасад над canonical Candle.
def get_klines(
self,
symbol: str | None = None,
*,
interval: str = "1m",
limit: int = 200,
price_type: str = "bid",
) -> KlineBatch:
symbol_to_use = symbol or self.settings.default_symbol
if not self.settings.exchange_enabled:
raise ExchangeError(
"Klines are not available in mock exchange mode."
)
if limit <= 0:
normalized_limit = 200
elif limit > 200:
normalized_limit = 200
else:
normalized_limit = limit
normalized_price_type = price_type.strip().lower()
if normalized_price_type not in {"bid", "ask"}:
normalized_price_type = "bid"
canonical_candles = self.get_candles(
symbol=symbol_to_use,
interval=interval,
limit=normalized_limit,
price_type=normalized_price_type,
)
candles = [ candles = [
self._kline_from_candle(candle) self._kline_from_candle(candle)
for candle in canonical_candles[-limit:] for candle in canonical_candles[-normalized_limit:]
] ]
normalized_symbol = (
canonical_candles[0].symbol
if canonical_candles
else normalize_symbol(symbol_to_use)
)
return KlineBatch( return KlineBatch(
symbol=validation.normalized_symbol, symbol=normalized_symbol,
interval=interval, interval=interval,
candles=candles, candles=candles,
source=f"rest_klines:{normalized_price_type}", source=f"rest_klines:{normalized_price_type}",

View File

@@ -0,0 +1,391 @@
# app/tests/unit/integrations/exchange/test_service_candles.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
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 _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 _candle(
*,
symbol: str = "BTC/USD_LEVERAGE",
interval: str = "1m",
open_time_ms: int = 1_750_000_000_000,
) -> Candle:
return Candle(
symbol=symbol,
interval=interval,
open_time=datetime.fromtimestamp(
open_time_ms / 1000,
tz=timezone.utc,
),
open_price=Decimal("100.10"),
high_price=Decimal("110.20"),
low_price=Decimal("90.30"),
close_price=Decimal("105.40"),
volume=Decimal("12.50"),
source="rest_klines:bid",
)
def test_get_candles_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_candles()
assert result == ()
assert requested_symbols == ["ETH/USD_LEVERAGE"]
assert acquisition_calls == [
{
"symbol": "ETH/USD_LEVERAGE",
"interval": "1m",
"limit": 200,
"price_type": "bid",
}
]
@pytest.mark.parametrize("limit", [0, -1, -100])
def test_get_candles_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_candles(limit=limit)
assert captured[0]["limit"] == 200
def test_get_candles_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_candles(limit=500)
assert captured[0]["limit"] == 200
@pytest.mark.parametrize("interval", ["1m", "5m", "15m", "1h"])
def test_get_candles_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 (),
)
service.get_candles(interval=interval)
assert captured[0]["interval"] == interval
def test_get_candles_rejects_unsupported_interval() -> None:
service = _service()
with pytest.raises(
ExchangeError,
match="Unsupported kline interval",
):
service.get_candles(interval="4h")
@pytest.mark.parametrize(
("price_type", "expected"),
[
("bid", "bid"),
("ask", "ask"),
(" BID ", "bid"),
(" AsK ", "ask"),
("unknown", "bid"),
("", "bid"),
],
)
def test_get_candles_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 (),
)
service.get_candles(price_type=price_type)
assert captured[0]["price_type"] == expected
def test_get_candles_rejects_mock_mode() -> None:
service = _service(exchange_enabled=False)
with pytest.raises(
ExchangeError,
match="Candles are not available in mock exchange mode",
):
service.get_candles()
def test_get_candles_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_candles("UNKNOWN")
def test_get_candles_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_candles(
" 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_candles_preserves_result_identity() -> None:
service = _service()
candles = (_candle(),)
_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_candles()
assert result is candles
def test_get_candles_preserves_order() -> None:
service = _service()
first = _candle(open_time_ms=1000)
second = _candle(open_time_ms=2000)
candles = (second, first)
_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_candles()
assert result is candles
assert result == (second, first)
def test_get_candles_preserves_empty_tuple_identity() -> None:
service = _service()
candles: tuple[Candle, ...] = ()
_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_candles()
assert result is candles
def test_get_candles_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_candles(
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",
},
}
]

View File

@@ -0,0 +1,569 @@
# Build 049 — Предоставление канонических Candle через публичный API ExchangeService
## Статус
**Завершён**
---
## Цель Build
Предоставить существующим потребителям публичный канонический API получения рыночных свечей:
```text
ExchangeService.get_candles()
tuple[Candle, ...]
```
без прямого доступа слоя Trading к деталям сборки `Market Data Acquisition` и без удаления действующего compatibility-контракта:
```text
ExchangeService.get_klines()
KlineBatch
```
Build 049 подготавливает orchestration-компоненты `Market Analysis` к последующему поэтапному переключению с legacy-моделей `Kline` и `KlineBatch` на каноническую модель `Candle`.
---
## Исходное состояние
До Build 049 `ExchangeService` уже использовал канонический Candles pipeline:
```text
DzengiCandlesDocumentSource
DzengiCandlesDocumentHandler
CandlesFeed
CandlesFeedRegistry
CandlesAcquisitionService
tuple[Candle, ...]
```
Однако публично сервис предоставлял только legacy-метод:
```text
ExchangeService.get_klines()
```
который:
1. выполнял проверку параметров;
2. вызывал канонический Acquisition pipeline;
3. преобразовывал `Candle` в `Kline`;
4. возвращал `KlineBatch`.
Из-за отсутствия публичного метода `get_candles()` непосредственное переключение `Market Analysis` на канонические модели потребовало бы либо использования внутреннего helper, либо дублирования composition chain в слое Trading.
---
## Объём изменений
В Build 049 изменены:
```text
src/integrations/exchange/service.py
tests/unit/integrations/exchange/test_service_candles.py
docs/migrations/build_049.md
```
Существующий файл:
```text
tests/unit/integrations/exchange/test_service_klines.py
```
не изменялся. Он использован как regression-контракт сохранения поведения legacy API.
Build 049 не изменяет:
```text
src/trading/market_analysis/service.py
src/trading/market_analysis/htf.py
src/market_data/acquisition/
src/integrations/exchange/models.py
```
Build не удаляет:
```text
ExchangeService.get_klines()
Kline
KlineBatch
_kline_from_candle()
```
---
## Новый публичный метод get_candles()
В `ExchangeService` добавлен метод:
```python
def get_candles(
self,
symbol: str | None = None,
*,
interval: str = "1m",
limit: int = 200,
price_type: str = "bid",
) -> tuple[Candle, ...]:
...
```
Метод возвращает immutable-набор канонических моделей:
```text
tuple[Candle, ...]
```
без преобразования в legacy-модели, без копирования и без сортировки результата.
---
## Ответственность get_candles()
`get_candles()` выполняет:
1. выбор `default_symbol`, если symbol не передан;
2. нормализацию `limit`;
3. проверку поддерживаемого interval;
4. нормализацию `price_type`;
5. проверку mock mode;
6. валидацию symbol;
7. вызов `_load_candles_via_acquisition()`;
8. логирование ошибок Acquisition pipeline;
9. преобразование ошибки в `ExchangeError`.
Нормализация limit:
```text
limit <= 0 → 200
limit > 200 → 200
иначе → исходное значение
```
Поддерживаемые intervals:
```text
1m
5m
15m
1h
```
Поддерживаемые значения `price_type`:
```text
bid
ask
```
Неизвестное значение нормализуется в:
```text
bid
```
---
## Канонический путь получения данных
После Build 049 публичный канонический путь выглядит так:
```text
ExchangeService.get_candles()
_load_candles_via_acquisition()
DzengiCandlesDocumentSource
DzengiCandlesDocumentHandler
CandlesFeed
CandlesFeedRegistry
CandlesAcquisitionService
tuple[Candle, ...]
```
Слой Trading получает стабильную публичную точку доступа к каноническим свечам и не знает о:
```text
DzengiCandlesDocumentSource
DzengiCandlesDocumentHandler
CandlesFeedRegistry
CandlesAcquisitionService composition
```
---
## Изменение get_klines()
`ExchangeService.get_klines()` сохранён как legacy compatibility-фасад.
После Build 049 его цепочка:
```text
ExchangeService.get_klines()
ExchangeService.get_candles()
tuple[Candle, ...]
_kline_from_candle()
list[Kline]
KlineBatch
```
`get_klines()` больше не вызывает `_load_candles_via_acquisition()` напрямую.
Таким образом, логика получения, проверки и обработки Acquisition-ошибок сосредоточена в одном публичном каноническом методе:
```text
get_candles()
```
---
## Сохранение legacy-контракта
Для `get_klines()` сохранены:
- публичная сигнатура;
- возвращаемый тип `KlineBatch`;
- нормализация limit;
- нормализация `price_type`;
- compatibility mapping `Candle → Kline`;
- источник вида `rest_klines:<price_type>`;
- прежнее сообщение mock mode:
```text
Klines are not available in mock exchange mode.
```
Для нового канонического API используется отдельное сообщение:
```text
Candles are not available in mock exchange mode.
```
Это сохраняет прежний контракт legacy-метода и одновременно вводит семантически корректный контракт нового API.
---
## Исключение повторной валидации symbol
При пустом наборе канонических свечей `get_klines()` не выполняет повторный вызов:
```text
validate_symbol()
```
Для заполнения symbol в пустом `KlineBatch` используется локальная нормализация:
```text
normalize_symbol(symbol_to_use)
```
Это сохраняет:
- один validation flow;
- один вызов `validate_symbol()`;
- корректный symbol пустого `KlineBatch`;
- отсутствие повторной загрузки Instrument Reference Data.
---
## Новый тестовый файл
Добавлен:
```text
tests/unit/integrations/exchange/test_service_candles.py
```
Тесты проверяют:
- использование `default_symbol`;
- нормализацию неположительного limit;
- ограничение limit значением 200;
- допустимые intervals;
- отклонение неподдерживаемого interval;
- нормализацию `price_type`;
- ошибку mock mode;
- ошибку invalid symbol;
- точные параметры вызова Acquisition helper;
- сохранение identity возвращённого tuple;
- сохранение порядка свечей;
- сохранение identity пустого tuple;
- логирование Acquisition-ошибок;
- преобразование ошибки в `ExchangeError`;
- сохранение исходного исключения как `__cause__`.
---
## Результаты targeted-тестов get_candles()
Выполнена команда:
```bash
python -m pytest -q tests/unit/integrations/exchange/test_service_candles.py
```
Результат:
```text
23 passed in 0.13s
```
---
## Совместная проверка canonical и legacy API
Выполнена команда:
```bash
python -m pytest -q tests/unit/integrations/exchange/test_service_candles.py tests/unit/integrations/exchange/test_service_klines.py
```
Результат:
```text
49 passed in 0.13s
```
Это подтверждает одновременно:
```text
get_candles() → новый канонический контракт
get_klines() → сохранённый legacy compatibility-контракт
```
---
## Полный regression suite
Выполнена команда:
```bash
python -m pytest -q
```
Результат:
```text
789 passed in 2.80s
```
Регрессий не обнаружено.
---
## Архитектурная проверка get_candles()
Выполнена команда:
```bash
grep -RIn --exclude-dir="__pycache__" --exclude="*.pyc" "def get_candles\|\.get_candles(" src tests
```
Результат подтверждает:
- определение `get_candles()` находится в `ExchangeService`;
- `get_klines()` вызывает `get_candles()`;
- специализированные вызовы находятся в `test_service_candles.py`;
- production-компоненты `Market Analysis` пока не переключены.
---
## Контроль legacy-потребителей
Выполнена команда:
```bash
grep -RIn --exclude-dir="__pycache__" --exclude="*.pyc" "\.get_klines(" src/trading/market_analysis tests/unit/integrations/exchange
```
Production-вызовы сохранены в:
```text
src/trading/market_analysis/service.py
src/trading/market_analysis/htf.py
```
Всего остаётся три production-вызова:
```text
MarketAnalysisService — один вызов
HTF orchestration — два вызова
```
Это ожидаемое переходное состояние.
---
## Контроль вызова Acquisition helper
Выполнена команда:
```bash
grep -RIn --exclude-dir="__pycache__" --exclude="*.pyc" "_load_candles_via_acquisition(" src/integrations/exchange/service.py tests/unit/integrations/exchange
```
В production обнаружены:
```text
один вызов внутри get_candles()
одно определение _load_candles_via_acquisition()
```
`get_klines()` больше не обращается к Acquisition helper напрямую.
---
## Контроль composition и compatibility bridge
Выполнена команда:
```bash
grep -nE "def get_candles|def get_klines|self\.get_candles|_load_candles_via_acquisition|_kline_from_candle" src/integrations/exchange/service.py
```
Результат подтверждает ожидаемую последовательность:
```text
get_candles()
_load_candles_via_acquisition()
get_klines()
self.get_candles()
_kline_from_candle()
```
---
## Проверка форматирования
Выполнена команда:
```bash
git diff --check
```
Вывод отсутствует.
Whitespace-ошибок не обнаружено.
---
## Архитектурный результат
После Build 049 доступны два явно разделённых публичных контракта.
Канонический API:
```text
ExchangeService.get_candles()
tuple[Candle, ...]
```
Legacy compatibility API:
```text
ExchangeService.get_klines()
get_candles()
Candle → Kline
KlineBatch
```
Получение данных и обработка ошибок больше не дублируются между двумя публичными методами.
---
## Что намеренно не выполнено
Build 049 намеренно не включает:
- переключение `MarketAnalysisService` на `get_candles()`;
- переключение HTF orchestration на `get_candles()`;
- удаление `get_klines()`;
- удаление `Kline`;
- удаление `KlineBatch`;
- удаление `_kline_from_candle()`;
- изменение торговой логики;
- изменение алгоритмов Market Analysis;
- изменение `Market Data Acquisition`;
- изменение структуры каталогов;
- удаление рабочего compatibility-кода.
---
## Критерии завершения
Build 049 считается завершённым, поскольку:
- добавлен публичный `ExchangeService.get_candles()`;
- метод возвращает канонический `tuple[Candle, ...]`;
- `get_klines()` использует `get_candles()` как единственный источник свечей;
- прямой вызов Acquisition helper из `get_klines()` устранён;
- legacy-контракт `KlineBatch` сохранён;
- повторная symbol validation устранена;
- новый canonical API покрыт специализированными тестами;
- legacy API проходит существующие regression-тесты;
- полный suite проходит;
- архитектурные grep соответствуют переходному состоянию;
- `git diff --check` чистый.
---
## Итог
**Build 049 завершён успешно.**
Публичный канонический API:
```text
ExchangeService.get_candles()
```
готов для поэтапного переключения orchestration-компонентов `Market Analysis`.
Текущее состояние:
```text
Targeted canonical + legacy tests: 49 passed
Full test suite: 789 passed
git diff --check: clean
```
Следующий безопасный этап — отдельное переключение непосредственных production-потребителей с:
```text
ExchangeService.get_klines()
```
на:
```text
ExchangeService.get_candles()
```
без одновременного удаления legacy compatibility-контракта.