build 041: remove legacy websocket quote parsing
This commit is contained in:
@@ -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 = {}
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user