Summary
WebSocketStreamBase keeps its stream-to-connection registry in two module-level objects (global_stream_connections, global_user_stream_connections in common/src/binance_common/websocket.py). Every client in the process shares them, keyed only by stream name.
Spot and USDⓈ-M Futures use the same stream names (btcusdt@aggTrade, btcusdt@depth, ...), and so do mainnet and testnet. When two stream clients in one process subscribe to the same name:
- the second client never sends
SUBSCRIBE, because the name is already registered;
- its callback is attached to the first client's connection, so it receives the other venue's data, with no error.
Reproduce
Self-contained: one local aiohttp server with two endpoints standing in for two venues; each tags its messages with its name. binance-common==4.5.1 (same as master at 43f19a9), Python 3.13.
import asyncio, json, logging
from aiohttp import web
from binance_common.configuration import ConfigurationWebSocketStreams
from binance_common.websocket import WebSocketStreamBase, RequestStream
logging.disable(logging.CRITICAL)
subs = {"spot": [], "futures": []}
def mk(name):
async def handler(request):
ws = web.WebSocketResponse(); await ws.prepare(request)
async for msg in ws:
d = json.loads(msg.data)
if d["method"] == "SUBSCRIBE":
subs[name] += d["params"]
await ws.send_json({"id": d["id"], "result": None})
for s in d["params"]:
await ws.send_json({"stream": s, "data": {"venue": name}})
return ws
return handler
async def main():
app = web.Application()
app.router.add_get("/spot/stream", mk("spot")); app.router.add_get("/fut/stream", mk("futures"))
r = web.AppRunner(app); await r.setup(); await web.TCPSite(r, "127.0.0.1", 18766).start()
spot = WebSocketStreamBase(ConfigurationWebSocketStreams(stream_url="ws://127.0.0.1:18766/spot/stream"))
fut = WebSocketStreamBase(ConfigurationWebSocketStreams(stream_url="ws://127.0.0.1:18766/fut/stream"))
await spot.create_connection(); await fut.create_connection()
got = {"spot_cb": [], "fut_cb": []}
h1 = await RequestStream(spot, "btcusdt@aggTrade"); h1.on("message", lambda d: got["spot_cb"].append(d))
h2 = await RequestStream(fut, "btcusdt@aggTrade"); h2.on("message", lambda d: got["fut_cb"].append(d))
await asyncio.sleep(0.5)
print("SUBSCRIBE frames received by servers:", subs)
print("callback payloads:", got)
await spot.close_connection(); await fut.close_connection(); await r.cleanup()
asyncio.run(main())
Output:
SUBSCRIBE frames received by servers: {'spot': ['btcusdt@aggTrade'], 'futures': []}
callback payloads: {'spot_cb': [{'stream': 'btcusdt@aggTrade', 'data': {'venue': 'spot'}}], 'fut_cb': [{'stream': 'btcusdt@aggTrade', 'data': {'venue': 'spot'}}]}
Expected
Each server receives its own SUBSCRIBE, and fut_cb receives {'venue': 'futures'}.
Impact
A process that follows the same symbol on Spot and Futures (basis or arbitrage monitoring, for example) gets Spot trades delivered to the Futures handler. Nothing is logged, and the payload shapes are close enough that it is easy to miss.
Suggested fix
Make both maps instance attributes created in the client's __init__ and use them wherever the globals are referenced today (subscribe filter, on, unsubscribe and the reconnect cleanup). Streams would then be deduplicated per client, which is what the pooling logic needs.
Summary
WebSocketStreamBasekeeps its stream-to-connection registry in two module-level objects (global_stream_connections,global_user_stream_connectionsincommon/src/binance_common/websocket.py). Every client in the process shares them, keyed only by stream name.Spot and USDⓈ-M Futures use the same stream names (
btcusdt@aggTrade,btcusdt@depth, ...), and so do mainnet and testnet. When two stream clients in one process subscribe to the same name:SUBSCRIBE, because the name is already registered;Reproduce
Self-contained: one local aiohttp server with two endpoints standing in for two venues; each tags its messages with its name.
binance-common==4.5.1(same as master at 43f19a9), Python 3.13.Output:
Expected
Each server receives its own
SUBSCRIBE, andfut_cbreceives{'venue': 'futures'}.Impact
A process that follows the same symbol on Spot and Futures (basis or arbitrage monitoring, for example) gets Spot trades delivered to the Futures handler. Nothing is logged, and the payload shapes are close enough that it is easy to miss.
Suggested fix
Make both maps instance attributes created in the client's
__init__and use them wherever the globals are referenced today (subscribe filter,on, unsubscribe and the reconnect cleanup). Streams would then be deduplicated per client, which is what the pooling logic needs.