Skip to content

[Bug] WebSocket Streams: stream registry is global, two clients with the same stream name get each other's data #570

Description

@Amadeus-22

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions