openlimits_binance/client/
stream.rs1use 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
33pub 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}