1use ring::hmac;
2use std::net::TcpStream;
3use tungstenite::handshake::client::Response;
4use tungstenite::stream::MaybeTlsStream;
5use tungstenite::{connect, Message, WebSocket};
6
7use std::sync::mpsc::{self, channel};
8use url::Url;
9
10use crate::errors::*;
11use crate::events::*;
12use chrono::Local;
13
14static SUBSCRIBED: &'static str = "subscribed";
15static INFO: &'static str = "info";
16static ERROR: &'static str = "error";
17static PONG: &'static str = "pong";
18static WEBSOCKET_URL: &'static str = "wss://ftx.com/ws/";
19
20pub trait EventHandler {
21 fn on_connect(&mut self, event: NotificationEvent);
22 fn on_auth(&mut self, event: NotificationEvent);
23 fn on_subscribed(&mut self, event: NotificationEvent);
24 fn on_data_event(&mut self, event: DataEvent);
25 fn on_error(&mut self, message: Error);
26}
27
28#[derive(Debug)]
29enum WsMessage {
30 Close,
31 Text(String),
32}
33
34pub struct WebSockets {
35 api_key: String,
36 secret_key: String,
37 socket: Option<(WebSocket<MaybeTlsStream<TcpStream>>, Response)>,
38 sender: Sender,
39 rx: mpsc::Receiver<WsMessage>,
40 event_handler: Option<Box<dyn EventHandler>>,
41 login_status: bool,
42}
43
44impl WebSockets {
45 pub fn new(api_key: Option<String>, secret_key: Option<String>) -> WebSockets {
46 let (tx, rx) = channel::<WsMessage>();
47 let sender = Sender { tx: tx };
48
49 WebSockets {
50 api_key: api_key.unwrap_or("".into()),
51 secret_key: secret_key.unwrap_or("".into()),
52 socket: None,
53 sender: sender,
54 rx: rx,
55 event_handler: None,
56 login_status: false,
57 }
58 }
59
60 pub fn connect(&mut self) -> Result<()> {
61 let wss: String = format!("{}", WEBSOCKET_URL);
62 let url = Url::parse(&wss)?;
63
64 match connect(url) {
65 Ok(answer) => {
66 self.socket = Some(answer);
67 Ok(())
68 }
69 Err(e) => {
70 bail!(format!("Error during handshake {}", e))
71 }
72 }
73 }
74
75 pub fn disconnect(&mut self) -> Result<()> {
76 if let Some(ref mut socket) = self.socket {
77 socket.0.close(None)?;
78 return Ok(());
79 }
80 bail!("Not able to close the connection");
81 }
82
83 pub fn add_event_handler<H>(&mut self, handler: H)
84 where
85 H: EventHandler + 'static,
86 {
87 self.event_handler = Some(Box::new(handler));
88 }
89
90 pub fn ping(&mut self) {
91 let msg = json!({
92 "op": "ping",
93 });
94 if let Err(e) = self.sender.send(&msg.to_string()) {
95 println!("{:?}", e);
96 }
97 }
98
99 pub fn subscribe_ticker<S>(&mut self, symbol: S)
100 where
101 S: Into<String>,
102 {
103 let msg = json!({
104 "op": "subscribe",
105 "channel": "ticker",
106 "market": symbol.into(),
107 });
108 if let Err(e) = self.sender.send(&msg.to_string()) {
109 println!("{:?}", e);
110 }
111 }
112
113 pub fn subscribe_trades<S>(&mut self, symbol: S)
114 where
115 S: Into<String>,
116 {
117 let msg = json!({
118 "op": "subscribe",
119 "channel": "trades",
120 "market": symbol.into(),
121 });
122 if let Err(e) = self.sender.send(&msg.to_string()) {
123 println!("{:?}", e);
124 }
125 }
126
127 pub fn subscribe_orderbook_grouped<S>(&mut self, symbol: S, group: i64)
128 where
129 S: Into<String>,
130 {
131 let msg = json!({
132 "op": "subscribe",
133 "channel": "orderbookGrouped",
134 "market": symbol.into(),
135 "grouping": group,
136 });
137 if let Err(e) = self.sender.send(&msg.to_string()) {
138 println!("{:?}", e);
139 }
140 }
141
142 pub fn subscribe_orderbook<S>(&mut self, symbol: S)
143 where
144 S: Into<String>,
145 {
146 let msg = json!({
147 "op": "subscribe",
148 "channel": "orderbook",
149 "market": symbol.into(),
150 });
151 if let Err(e) = self.sender.send(&msg.to_string()) {
152 println!("{:?}", e);
153 }
154 }
155
156 pub fn unsubscribe(&mut self, channel: String, symbol: Option<String>, grouping: Option<i64>) {
157 let mut msg = json!({
158 "op": "unsubscribe",
159 "channel": channel,
160 });
161 if let Some(ref market) = symbol {
162 msg = json!({
163 "op": "unsubscribe",
164 "channel": channel,
165 "market": market,
166 });
167 }
168 if let Some(g) = grouping {
169 msg = json!({
170 "op": "unsubscribe",
171 "channel": channel,
172 "market": &symbol.unwrap(),
173 "grouping": g,
174 });
175 }
176 if let Err(e) = self.sender.send(&msg.to_string()) {
177 println!("{:?}", e);
178 }
179 }
180
181 pub fn login(&mut self) {
182 let ts = Local::now().timestamp() * 1000;
183 let signature_payload = format!("{}websocket_login", ts);
184 let signed_key = hmac::Key::new(hmac::HMAC_SHA256, self.secret_key.as_bytes());
185 let signature = hex::encode(hmac::sign(&signed_key, signature_payload.as_bytes()).as_ref());
186
187 let msg = json!({
188 "op": "login",
189 "args": {
190 "key": self.api_key,
191 "sign": signature,
192 "time": ts,
193 },
194 });
195 if let Err(e) = self.sender.send(&msg.to_string()) {
196 println!("{:?}", e);
197 }
198 self.login_status = true;
199 }
200
201 pub fn subscribe_fills(&mut self) {
202 if !self.login_status {
203 self.login();
204 }
205
206 let msg = json!({
207 "op": "subscribe",
208 "channel": "fills",
209 });
210 if let Err(e) = self.sender.send(&msg.to_string()) {
211 println!("{:?}", e);
212 }
213 }
214
215 pub fn subscribe_orders(&mut self) {
216 if !self.login_status {
217 self.login();
218 }
219
220 let msg = json!({
221 "op": "subscribe",
222 "channel": "orders",
223 });
224 if let Err(e) = self.sender.send(&msg.to_string()) {
225 println!("{:?}", e);
226 }
227 }
228
229 pub fn subscribe_ftxpay(&mut self) {
230 if !self.login_status {
231 self.login();
232 }
233
234 let msg = json!({
235 "op": "subscribe",
236 "channel": "ftxpay",
237 });
238 if let Err(e) = self.sender.send(&msg.to_string()) {
239 println!("{:?}", e);
240 }
241 }
242
243 pub fn event_loop(&mut self) -> Result<()> {
244 let mut ping_flag = 0;
245 loop {
246 if let Some(ref mut socket) = self.socket {
247 loop {
248 match self.rx.try_recv() {
249 Ok(msg) => match msg {
250 WsMessage::Text(text) => {
251 socket.0.write_message(Message::Text(text))?;
252 }
253 WsMessage::Close => {
254 return socket.0.close(None).map_err(|e| e.into());
255 }
256 },
257 Err(mpsc::TryRecvError::Disconnected) => {
258 bail!("Disconnected")
259 }
260 Err(mpsc::TryRecvError::Empty) => break,
261 }
262 }
263
264 let message = socket.0.read_message()?;
265
266 match message {
267 Message::Text(text) => {
268 if let Some(ref mut h) = self.event_handler {
269 if text.find(INFO) != None {
270 println!("INFO: {:?}", text);
271 } else if text.find(SUBSCRIBED) != None {
272 let event: NotificationEvent = serde_json::from_str(&text)?;
273 h.on_subscribed(event);
274 } else if text.find(ERROR) != None {
275 println!("ERROR: {:?}", text);
276 } else if text.find(PONG) != None {
277 let event: NotificationEvent = serde_json::from_str(&text)?;
278 h.on_subscribed(event);
279 } else {
280 let event: DataEvent = serde_json::from_str(&text)?;
282 h.on_data_event(event);
283 }
284 }
285 }
286 Message::Binary(_) => {}
287 Message::Ping(_) | Message::Pong(_) => {}
288 Message::Close(e) => {
289 bail!(format!("Disconnected {:?}", e));
290 }
291 }
292
293 ping_flag = ping_flag + 1;
294 if ping_flag >= 20 {
295 ping_flag = 0;
296 self.ping();
297 }
298 }
299 }
300 }
301}
302
303#[derive(Clone)]
304pub struct Sender {
305 tx: mpsc::Sender<WsMessage>,
306}
307
308impl Sender {
309 pub fn send(&self, raw: &str) -> Result<()> {
310 self.tx
311 .send(WsMessage::Text(raw.to_string()))
312 .map_err(|e| Error::with_chain(e, "not able to send a message"))?;
313 Ok(())
314 }
315
316 pub fn shutdown(&self) -> Result<()> {
317 self.tx
318 .send(WsMessage::Close)
319 .map_err(|e| Error::with_chain(e, "Error during shutdown"))?;
320 Ok(())
321 }
322}