diff --git a/app/src/integrations/exchange/market_data_runner.py b/app/src/integrations/exchange/market_data_runner.py index e0d42e1..8e5204d 100644 --- a/app/src/integrations/exchange/market_data_runner.py +++ b/app/src/integrations/exchange/market_data_runner.py @@ -8,8 +8,7 @@ import traceback from dataclasses import dataclass from typing import Callable -from src.core.numbers import safe_float -from src.core.types import JsonDict, NumericLike +from src.core.types import JsonDict from src.integrations.exchange.market_cache import MarketPriceCache from src.integrations.exchange.service import ExchangeService from src.integrations.exchange.ws_client import ExchangeWebSocketClient @@ -490,36 +489,6 @@ class MarketDataRunner: def _ws_symbol(cls, symbol: str) -> str: return cls._cache_symbol(symbol) - @classmethod - def _extract_best_price( - cls, - payload: JsonDict, - side_key: str, - ) -> float | None: - data = cls._extract_depth_payload(payload) - - values = data.get(side_key) - - if not isinstance(values, list) or not values: - return None - - first = values[0] - - if isinstance(first, list) and first: - return cls._positive_float(first[0]) - - if isinstance(first, dict): - raw_price = ( - first.get("price") - or first.get("p") - or first.get("bidPrice") - or first.get("askPrice") - ) - - return cls._positive_float(raw_price) - - return None - @classmethod def _extract_depth_payload(cls, payload: JsonDict) -> JsonDict: data: object = payload @@ -538,15 +507,6 @@ class MarketDataRunner: return payload - @classmethod - def _positive_float(cls, value: NumericLike | None) -> float | None: - number = safe_float(value) - - if number is None or number <= 0: - return None - - return number - @classmethod def _safe_payload_preview(cls, payload: JsonDict) -> JsonDict: preview: JsonDict = {} diff --git a/app/src/integrations/exchange/market_stream.py b/app/src/integrations/exchange/market_stream.py index 2880001..be578fc 100644 --- a/app/src/integrations/exchange/market_stream.py +++ b/app/src/integrations/exchange/market_stream.py @@ -3,12 +3,7 @@ from __future__ import annotations import asyncio -from datetime import datetime -from zoneinfo import ZoneInfo - from src.core.config import load_settings -from src.core.numbers import safe_float -from src.core.types import JsonDict, NumericLike from src.integrations.exchange.market_cache import MarketPriceCache from src.integrations.exchange.service import ExchangeService from src.integrations.exchange.ws_client import ExchangeWebSocketClient @@ -21,117 +16,6 @@ from src.market_data.acquisition.exceptions import ( from src.trading.journal.service import JournalService -# безопасно форматирует timestamp биржи в локальное время -def _format_timestamp(raw_timestamp: NumericLike | None) -> str | None: - timestamp = safe_float(raw_timestamp) - - if timestamp is None: - return None - - try: - settings = load_settings() - - dt_utc = datetime.fromtimestamp( - int(timestamp) / 1000, - tz=ZoneInfo("UTC"), - ) - - return dt_utc.astimezone( - ZoneInfo(settings.tz), - ).strftime("%d.%m.%Y %H:%M:%S") - - except Exception: - return None - - -# достаёт внутренний payload из websocket-сообщения -def _payload_from_message(payload: JsonDict) -> JsonDict | None: - event = payload.get("Payload") or payload.get("payload") - - if isinstance(event, dict) and "Payload" in event: - event = event.get("Payload") - - if not isinstance(event, dict): - return None - - return dict(event) - - -# извлекает best bid / best ask из формата depth -def _extract_depth_prices(event: JsonDict) -> tuple[float | None, float | None]: - bids = event.get("bids") - asks = event.get("asks") - - bid_price = _extract_first_price(bids) - ask_price = _extract_first_price(asks) - - return bid_price, ask_price - - -# извлекает первую цену из списка стакана -def _extract_first_price(value: object) -> float | None: - if not isinstance(value, list) or not value: - return None - - first = value[0] - - if isinstance(first, list) and first: - return _positive_float(first[0]) - - if isinstance(first, dict): - return _positive_float( - first.get("price") - or first.get("p") - or first.get("bidPrice") - or first.get("askPrice") - ) - - return None - - -# безопасно приводит число к float и отсекает нулевые/отрицательные цены -def _positive_float(value: NumericLike | None) -> float | None: - number = safe_float(value) - - if number is None or number <= 0: - return None - - return number - - -# нормализует websocket-сообщение рынка в единый формат для MarketPriceCache -def _extract_market_event(payload: JsonDict) -> JsonDict | None: - event = _payload_from_message(payload) - - if event is None: - return None - - symbol = ( - event.get("symbolName") - or event.get("symbol") - or payload.get("symbol") - ) - - bid_price = _positive_float(event.get("bid")) - ask_price = _positive_float(event.get("ofr") or event.get("ask")) - - if bid_price is None or ask_price is None: - bid_price, ask_price = _extract_depth_prices(event) - - if symbol is None or bid_price is None or ask_price is None: - return None - - price = (bid_price + ask_price) / 2 - - return { - "symbol": str(symbol).upper(), - "price": price, - "bid_price": bid_price, - "ask_price": ask_price, - "updated_at": _format_timestamp(event.get("timestamp")), - } - - # запускает постоянный websocket-поток рынка и обновляет MarketPriceCache async def start_market_stream() -> None: settings = load_settings() diff --git a/app/tests/unit/integrations/exchange/test_market_data_runner.py b/app/tests/unit/integrations/exchange/test_market_data_runner.py index fdd97dd..9e954e4 100644 --- a/app/tests/unit/integrations/exchange/test_market_data_runner.py +++ b/app/tests/unit/integrations/exchange/test_market_data_runner.py @@ -248,13 +248,6 @@ def test_run_websocket_uses_canonical_quote_adapter( monkeypatch.setattr(runner_module, "ExchangeWebSocketClient", Client) monkeypatch.setattr(runner_module, "DzengiWebSocketQuoteAdapter", Adapter) monkeypatch.setattr(runner_module, "MarketPriceCache", Cache) - monkeypatch.setattr( - MarketDataRunner, - "_extract_best_price", - lambda *args, **kwargs: (_ for _ in ()).throw( - AssertionError("Legacy parser must not be called.") - ), - ) asyncio.run( MarketDataRunner._run_websocket( diff --git a/app/tests/unit/integrations/exchange/test_market_stream.py b/app/tests/unit/integrations/exchange/test_market_stream.py index 13a6763..a25485c 100644 --- a/app/tests/unit/integrations/exchange/test_market_stream.py +++ b/app/tests/unit/integrations/exchange/test_market_stream.py @@ -380,13 +380,6 @@ def test_start_market_stream_maps_message_to_quote( monkeypatch.setattr(stream_module, "ExchangeWebSocketClient", Client) monkeypatch.setattr(stream_module, "DzengiWebSocketQuoteAdapter", Adapter) monkeypatch.setattr(stream_module, "MarketPriceCache", Cache) - monkeypatch.setattr( - stream_module, - "_extract_market_event", - lambda payload: (_ for _ in ()).throw( - AssertionError("Legacy parser must not be called.") - ), - ) monkeypatch.setattr(stream_module.asyncio, "sleep", _raise_stop_stream) with pytest.raises(StopStream):