"""币安公共 WebSocket:USDT 永续 bookTicker + 期权 bookTicker。只读。""" from __future__ import annotations import asyncio import json import logging from typing import Any from urllib.parse import urlparse import websockets from websockets.asyncio.client import ClientConnection from ..book_cache import BookCache from .parse import is_option_symbol, safe_float logger = logging.getLogger(__name__) class BinancePublicWs: def __init__( self, *, futures_ws_base: str, options_ws_base: str, cache: BookCache, proxy: str | None = None, ping_interval: float = 20.0, ) -> None: self.futures_ws_base = futures_ws_base.rstrip("/") self.options_ws_base = options_ws_base.rstrip("/") self.cache = cache self.proxy = (proxy or "").strip() or None self.ping_interval = ping_interval self._inst_ids: list[str] = [] self._tasks: list[asyncio.Task[None]] = [] self._stop = asyncio.Event() self._alive: set[str] = set() def set_instruments(self, inst_ids: list[str]) -> None: self._inst_ids = [i for i in inst_ids if i] def _split(self) -> tuple[list[str], list[str]]: perps: list[str] = [] opts: list[str] = [] for i in self._inst_ids: if is_option_symbol(i): opts.append(i) else: perps.append(i.upper()) return perps, opts def _mark_alive(self, kind: str, ok: bool) -> None: if ok: self._alive.add(kind) else: self._alive.discard(kind) # 任一腿连上即认为行情可用(期权腿偶发失败时仍可用 REST/永续) self.cache.set_connected(bool(self._alive)) async def start(self) -> None: if self._tasks and any(not t.done() for t in self._tasks): return self._stop.clear() await self._spawn() async def stop(self) -> None: self._stop.set() for t in self._tasks: t.cancel() for t in self._tasks: try: await t except asyncio.CancelledError: pass self._tasks = [] self._alive.clear() self.cache.set_connected(False) async def resubscribe(self, inst_ids: list[str]) -> None: self.set_instruments(inst_ids) await self.stop() self._stop.clear() await self._spawn() async def _spawn(self) -> None: perps, opts = self._split() self._tasks = [] self._alive.clear() if perps: url = self._combined_url(self.futures_ws_base, [f"{p.lower()}@bookTicker" for p in perps]) self._tasks.append( asyncio.create_task(self._run_forever(url, kind="futures"), name="bn-fapi-ws") ) if opts: # 期权 stream 名必须小写:eth-260726-1850-c@bookTicker streams = [f"{s.lower()}@bookTicker" for s in opts] url = self._combined_url(self.options_ws_base, streams) self._tasks.append( asyncio.create_task(self._run_forever(url, kind="options"), name="bn-eapi-ws") ) if not self._tasks: self.cache.set_connected(False) @staticmethod def _combined_url(base: str, streams: list[str]) -> str: # base like wss://fstream.binance.com/stream or .../eoptions/stream if "/stream" in base: root = base else: root = base.rstrip("/") + "/stream" return root + "?streams=" + "/".join(streams) async def _open_connection(self, url: str) -> ClientConnection: if not self.proxy: return await websockets.connect( url, ping_interval=None, max_size=2**22, open_timeout=20, ) from python_socks.async_.asyncio import Proxy parsed = urlparse(url) host = parsed.hostname or "fstream.binance.com" port = parsed.port or (443 if parsed.scheme == "wss" else 80) sock = await Proxy.from_url(self.proxy).connect(dest_host=host, dest_port=port) return await websockets.connect( url, sock=sock, server_hostname=host, ping_interval=None, max_size=2**22, open_timeout=20, ) async def _run_forever(self, url: str, *, kind: str) -> None: backoff = 1.0 while not self._stop.is_set(): try: async with await self._open_connection(url) as ws: self._mark_alive(kind, True) backoff = 1.0 logger.info("Binance %s WS connected: %s", kind, url[:160]) waiter = asyncio.create_task(self._stop.wait()) reader = asyncio.create_task(self._read_loop(ws)) pinger = asyncio.create_task(self._ping_loop(ws)) done, pending = await asyncio.wait( {waiter, reader, pinger}, return_when=asyncio.FIRST_COMPLETED, ) for t in pending: t.cancel() for t in done: exc = t.exception() if exc and not isinstance(exc, asyncio.CancelledError): raise exc except asyncio.CancelledError: raise except Exception as e: logger.warning("Binance %s WS disconnected: %s", kind, e) self._mark_alive(kind, False) try: await asyncio.wait_for(self._stop.wait(), timeout=backoff) break except asyncio.TimeoutError: backoff = min(backoff * 2, 30.0) self._mark_alive(kind, False) async def _ping_loop(self, ws: ClientConnection) -> None: while True: await asyncio.sleep(self.ping_interval) try: await ws.ping() except Exception: return async def _read_loop(self, ws: ClientConnection) -> None: try: async for raw in ws: if isinstance(raw, bytes): raw = raw.decode("utf-8", errors="ignore") try: msg = json.loads(raw) except json.JSONDecodeError: continue data = msg.get("data") if isinstance(msg, dict) and "stream" in msg else msg if isinstance(data, dict): self._handle_event(data) except websockets.exceptions.ConnectionClosed: return def _handle_event(self, data: dict[str, Any]) -> None: et = str(data.get("e") or "") sym = str(data.get("s") or "") if not sym: return ts = safe_float(data.get("E") or data.get("T")) ts_ms = int(ts) if ts is not None else None if et in ("bookTicker", "") or ("b" in data and "a" in data and "s" in data): bid = safe_float(data.get("b")) ask = safe_float(data.get("a")) bid_sz = safe_float(data.get("B")) ask_sz = safe_float(data.get("A")) if bid is not None or ask is not None: self.cache.upsert_top( sym, bid=bid, ask=ask, bid_sz=bid_sz, ask_sz=ask_sz, ts_ms=ts_ms, ) # 永续可用中间价近似 mark if not is_option_symbol(sym) and bid and ask: self.cache.set_mark_px(sym, (bid + ask) / 2.0, ts_ms=ts_ms)