build 047: switch ExchangeService klines to canonical candles feed
This commit is contained in:
@@ -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
|
||||
return None
|
||||
|
||||
Reference in New Issue
Block a user