Skip to main content

openlimits_binance/client/
stream.rs

1use std::{convert::TryFrom, fmt::Display};
2use std::sync::Mutex;
3use async_trait::async_trait;
4use futures::{SinkExt, stream::BoxStream, StreamExt};
5use serde::{de, Deserialize, Serialize};
6use serde_json::Value;
7use tokio::sync::mpsc::{unbounded_channel, UnboundedSender};
8use tokio_tungstenite::{connect_async, tungstenite::protocol::Message};
9use openlimits_exchange::errors::OpenLimitsError;
10use crate::{
11    BinanceParameters,
12    model::websocket::{BinanceSubscription, BinanceWebsocketMessage},
13};
14use openlimits_exchange::{
15    model::websocket::OpenLimitsWebSocketMessage,
16    model::websocket::Subscription,
17    model::websocket::WebSocketResponse,
18};
19use openlimits_exchange::traits::stream::{ExchangeStream, Subscriptions};
20use super::shared::Result;
21use openlimits_exchange::exchange::Environment;
22
23const WS_URL_PROD: &str = "wss://stream.binance.com:9443/stream";
24const WS_URL_SANDBOX: &str = "wss://testnet.binance.vision/stream";
25
26#[derive(Debug, Clone, Deserialize, Serialize)]
27#[serde(untagged)]
28enum Either<L, R> {
29    Left(L),
30    Right(R),
31}
32
33/// This struct is used for websocket communications with openlimits-binance openlimits-exchange
34pub struct BinanceWebsocket {
35    parameters: BinanceParameters,
36    disconnection_senders: Mutex<Vec<UnboundedSender<()>>>,
37}
38
39#[async_trait]
40impl ExchangeStream for BinanceWebsocket {
41    type InitParams = BinanceParameters;
42    type Subscription = BinanceSubscription;
43    type Response = BinanceWebsocketMessage;
44
45    async fn new(parameters: Self::InitParams) -> Result<Self> {
46        Ok(BinanceWebsocket {
47            parameters,
48            disconnection_senders: Default::default(),
49        })
50    }
51
52    async fn disconnect(&self) {
53        if let Ok(mut senders) = self.disconnection_senders.lock() {
54            for sender in senders.iter() {
55                sender.send(()).ok();
56            }
57            senders.clear();
58        }
59    }
60
61    async fn create_stream_specific(
62        &self,
63        subscriptions: Subscriptions<Self::Subscription>,
64    ) -> Result<BoxStream<'static, Result<Self::Response>>> {
65        let streams = subscriptions
66            .into_iter()
67            .map(|bs| bs.to_string())
68            .collect::<Vec<String>>()
69            .join("/");
70
71        let ws_url = match self.parameters.environment {
72            Environment::Sandbox => WS_URL_SANDBOX,
73            Environment::Production => WS_URL_PROD,
74        };
75        let endpoint = url::Url::parse(&format!("{}?streams={}", ws_url, streams.to_lowercase()))
76            .map_err(OpenLimitsError::UrlParserError)?;
77        let (ws_stream, _) = connect_async(endpoint).await?;
78
79        let (mut sink, stream) = ws_stream.split();
80        let (disconnection_sender, mut disconnection_receiver) = unbounded_channel();
81        tokio::spawn(async move {
82            if disconnection_receiver.recv().await.is_some() {
83                sink.close().await.ok();
84            }
85        });
86
87        if let Ok(mut senders) = self.disconnection_senders.lock() {
88            senders.push(disconnection_sender);
89        }
90
91        let s = stream.map(|message| match message {
92            Ok(msg) => parse_message(msg),
93            Err(_) => Err(OpenLimitsError::SocketError()),
94        });
95
96        Ok(s.boxed())
97    }
98}
99
100#[derive(Deserialize)]
101struct BinanceWebsocketStream {
102    #[serde(rename = "stream")]
103    pub name: String,
104    pub data: Value,
105}
106
107impl<'de> Deserialize<'de> for BinanceWebsocketMessage {
108    fn deserialize<D>(deserializer: D) -> core::result::Result<Self, D::Error>
109    where
110        D: serde::Deserializer<'de>,
111    {
112        let stream: BinanceWebsocketStream = BinanceWebsocketStream::deserialize(deserializer)?;
113
114        if stream.name.ends_with("@aggTrade") {
115            Ok(BinanceWebsocketMessage::AggregateTrade(
116                serde_json::from_value(stream.data).map_err(de::Error::custom)?,
117            ))
118        } else if stream.name.contains("@trade") {
119            Ok(BinanceWebsocketMessage::Trade(
120                serde_json::from_value(stream.data).map_err(de::Error::custom)?,
121            ))
122        } else if stream.name.contains("@kline_") {
123            Ok(BinanceWebsocketMessage::Candlestick(
124                serde_json::from_value(stream.data).map_err(de::Error::custom)?,
125            ))
126        } else if stream.name.contains("@ticker") {
127            Ok(BinanceWebsocketMessage::Ticker(
128                serde_json::from_value(stream.data).map_err(de::Error::custom)?,
129            ))
130        } else if stream.name.eq("!ticker@arr") {
131            Ok(BinanceWebsocketMessage::TickerAll(
132                serde_json::from_value(stream.data).map_err(de::Error::custom)?,
133            ))
134        } else if stream.name.ends_with("@miniTicker") {
135            Ok(BinanceWebsocketMessage::MiniTicker(
136                serde_json::from_value(stream.data).map_err(de::Error::custom)?,
137            ))
138        } else if stream.name.ends_with("!miniTicker@arr") {
139            Ok(BinanceWebsocketMessage::MiniTickerAll(
140                serde_json::from_value(stream.data).map_err(de::Error::custom)?,
141            ))
142        } else if stream.name.ends_with("@depth") {
143            Ok(BinanceWebsocketMessage::Depth(
144                serde_json::from_value(stream.data).map_err(de::Error::custom)?,
145            ))
146        } else if stream.name.contains("@depth") {
147            Ok(BinanceWebsocketMessage::OrderBook(
148                serde_json::from_value(stream.data).map_err(de::Error::custom)?,
149            ))
150        } else {
151            panic!("Not supported Subscription");
152        }
153    }
154}
155
156impl Display for BinanceSubscription {
157    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
158        match self {
159            BinanceSubscription::AggregateTrade(ref symbol) => write!(f, "{}@aggTrade", symbol),
160            BinanceSubscription::Candlestick(ref symbol, ref interval) => {
161                write!(f, "{}@kline_{}", symbol, interval)
162            }
163            BinanceSubscription::Depth(ref symbol, interval) => match interval {
164                None => write!(f, "{}@depth", symbol),
165                Some(i) => write!(f, "{}@depth@{}ms", symbol, i),
166            },
167            BinanceSubscription::MiniTicker(symbol) => write!(f, "{}@miniTicker", symbol),
168            BinanceSubscription::MiniTickerAll => write!(f, "!miniTicker@arr"),
169            BinanceSubscription::OrderBook(ref symbol, depth) => {
170                write!(f, "{}@depth{}", symbol, depth)
171            }
172            BinanceSubscription::Ticker(ref symbol) => write!(f, "{}@ticker", symbol),
173            BinanceSubscription::TickerAll => write!(f, "!ticker@arr"),
174            BinanceSubscription::Trade(ref symbol) => write!(f, "{}@trade", symbol),
175            BinanceSubscription::UserData(ref key) => write!(f, "{}", key),
176        }
177    }
178}
179
180impl From<Subscription> for BinanceSubscription {
181    fn from(subscription: Subscription) -> Self {
182        match subscription {
183            Subscription::OrderBookUpdates(symbol) => BinanceSubscription::Depth(crate::model::MarketPair::from(symbol).0, None),
184            Subscription::Trades(symbol) => BinanceSubscription::Trade(crate::model::MarketPair::from(symbol).0)
185        }
186    }
187}
188
189impl TryFrom<BinanceWebsocketMessage> for WebSocketResponse<BinanceWebsocketMessage> {
190    type Error = OpenLimitsError;
191
192    fn try_from(value: BinanceWebsocketMessage) -> Result<Self> {
193        match value {
194            BinanceWebsocketMessage::Depth(orderbook) => Ok(WebSocketResponse::Generic(
195                OpenLimitsWebSocketMessage::OrderBook(orderbook.into()),
196            )),
197            BinanceWebsocketMessage::Trade(trade) => Ok(WebSocketResponse::Generic(
198                OpenLimitsWebSocketMessage::Trades(trade.into()),
199            )),
200            BinanceWebsocketMessage::Ping => {
201                Ok(WebSocketResponse::Generic(OpenLimitsWebSocketMessage::Ping))
202            }
203            BinanceWebsocketMessage::Close => Err(OpenLimitsError::SocketError()),
204            _ => Ok(WebSocketResponse::Raw(value)),
205        }
206    }
207}
208
209fn parse_message(ws_message: Message) -> Result<BinanceWebsocketMessage> {
210    let msg = match ws_message {
211        Message::Text(m) => m,
212        Message::Binary(b) => return Ok(BinanceWebsocketMessage::Binary(b)),
213        Message::Pong(..) => return Ok(BinanceWebsocketMessage::Pong),
214        Message::Ping(..) => return Ok(BinanceWebsocketMessage::Ping),
215        Message::Close(..) => return Ok(BinanceWebsocketMessage::Close),
216    };
217
218    serde_json::from_str(&msg).map_err(OpenLimitsError::JsonError)
219}