nautilus_bitmex/websocket/
handler.rs1use std::sync::{
19 Arc,
20 atomic::{AtomicBool, Ordering},
21};
22
23use nautilus_network::{
24 RECONNECTED,
25 retry::{RetryManager, create_websocket_retry_manager},
26 websocket::{AuthTracker, SubscriptionState, WebSocketClient},
27};
28use tokio_tungstenite::tungstenite::Message;
29
30use super::{
31 enums::{BitmexWsAuthAction, BitmexWsOperation},
32 error::BitmexWsError,
33 messages::{BitmexHttpRequest, BitmexTableMessage, BitmexWsFrame, BitmexWsMessage},
34};
35
36#[derive(Debug)]
38pub enum HandlerCommand {
39 SetClient(WebSocketClient),
41 Disconnect,
43 Authenticate { payload: String },
45 Subscribe { topics: Vec<String> },
47 Unsubscribe { topics: Vec<String> },
49}
50
51pub(super) struct BitmexWsFeedHandler {
52 signal: Arc<AtomicBool>,
53 inner: Option<WebSocketClient>,
54 cmd_rx: tokio::sync::mpsc::UnboundedReceiver<HandlerCommand>,
55 raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
56 out_tx: tokio::sync::mpsc::UnboundedSender<BitmexWsMessage>,
57 auth_tracker: AuthTracker,
58 subscriptions: SubscriptionState,
59 retry_manager: RetryManager<BitmexWsError>,
60}
61
62impl BitmexWsFeedHandler {
63 pub(super) fn new(
65 signal: Arc<AtomicBool>,
66 cmd_rx: tokio::sync::mpsc::UnboundedReceiver<HandlerCommand>,
67 raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
68 out_tx: tokio::sync::mpsc::UnboundedSender<BitmexWsMessage>,
69 auth_tracker: AuthTracker,
70 subscriptions: SubscriptionState,
71 ) -> Self {
72 Self {
73 signal,
74 inner: None,
75 cmd_rx,
76 raw_rx,
77 out_tx,
78 auth_tracker,
79 subscriptions,
80 retry_manager: create_websocket_retry_manager(),
81 }
82 }
83
84 pub(super) fn is_stopped(&self) -> bool {
85 self.signal.load(Ordering::Relaxed)
86 }
87
88 pub(super) fn send(&self, msg: BitmexWsMessage) -> Result<(), ()> {
89 self.out_tx.send(msg).map_err(|_| ())
90 }
91
92 async fn send_with_retry(&self, payload: String) -> anyhow::Result<()> {
94 if let Some(client) = &self.inner {
95 self.retry_manager
96 .execute_with_retry(
97 "websocket_send",
98 || {
99 let payload = payload.clone();
100 async move {
101 client.send_text(payload, None).await.map_err(|e| {
102 BitmexWsError::ClientError(format!("Send failed: {e}"))
103 })
104 }
105 },
106 should_retry_bitmex_error,
107 |e| create_bitmex_timeout_error(e.to_string()),
108 )
109 .await
110 .map_err(|e| anyhow::anyhow!("{e}"))
111 } else {
112 Err(anyhow::anyhow!("No active WebSocket client"))
113 }
114 }
115
116 pub(super) async fn next(&mut self) -> Option<BitmexWsMessage> {
117 loop {
118 tokio::select! {
119 Some(cmd) = self.cmd_rx.recv() => {
120 match cmd {
121 HandlerCommand::SetClient(client) => {
122 log::debug!("WebSocketClient received by handler");
123 self.inner = Some(client);
124 }
125 HandlerCommand::Disconnect => {
126 log::debug!("Disconnect command received");
127
128 if let Some(client) = self.inner.take() {
129 client.disconnect().await;
130 }
131 }
132 HandlerCommand::Authenticate { payload } => {
133 log::debug!("Authenticate command received");
134
135 if let Err(e) = self.send_with_retry(payload).await {
136 log::error!("Failed to send authentication after retries: {e}");
137 }
138 }
139 HandlerCommand::Subscribe { topics } => {
140 for topic in topics {
141 log::debug!("Subscribing to topic: {topic}");
142 if let Err(e) = self.send_with_retry(topic.clone()).await {
143 log::error!("Failed to send subscription after retries: topic={topic}, error={e}");
144 }
145 }
146 }
147 HandlerCommand::Unsubscribe { topics } => {
148 for topic in topics {
149 log::debug!("Unsubscribing from topic: {topic}");
150 if let Err(e) = self.send_with_retry(topic.clone()).await {
151 log::error!("Failed to send unsubscription after retries: topic={topic}, error={e}");
152 }
153 }
154 }
155 }
156 }
157
158 () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {
159 if self.signal.load(std::sync::atomic::Ordering::Relaxed) {
160 log::debug!("Stop signal received during idle period");
161 return None;
162 }
163 }
164
165 msg = self.raw_rx.recv() => {
166 let msg = match msg {
167 Some(msg) => msg,
168 None => {
169 log::debug!("WebSocket stream closed");
170 return None;
171 }
172 };
173
174 if let Message::Ping(data) = &msg {
176 log::trace!("Received ping frame with {} bytes", data.len());
177
178 if let Some(client) = &self.inner
179 && let Err(e) = client.send_pong(data.to_vec()).await
180 {
181 log::warn!("Failed to send pong frame: {e}");
182 }
183 continue;
184 }
185
186 let event = match self.parse_raw_message(msg) {
187 Some(event) => event,
188 None => continue,
189 };
190
191 if self.signal.load(std::sync::atomic::Ordering::Relaxed) {
192 log::debug!("Stop signal received");
193 return None;
194 }
195
196 match event {
197 BitmexWsFrame::Reconnected => {
198 return Some(BitmexWsMessage::Reconnected);
199 }
200 BitmexWsFrame::Subscription {
201 success,
202 subscribe,
203 request,
204 error,
205 } => {
206 if let Some(msg) = self.handle_subscription_message(
207 success,
208 subscribe.as_ref(),
209 request.as_ref(),
210 error.as_deref(),
211 ) {
212 return Some(msg);
213 }
214 }
215 BitmexWsFrame::Table(table_msg) => {
216 return Some(BitmexWsMessage::Table(table_msg));
217 }
218 BitmexWsFrame::Welcome { .. } | BitmexWsFrame::Error { .. } => {}
219 }
220 }
221
222 else => {
224 log::debug!("Handler shutting down: stream ended or command channel closed");
225 return None;
226 }
227 }
228 }
229 }
230
231 fn parse_raw_message(&self, msg: Message) -> Option<BitmexWsFrame> {
232 match msg {
233 Message::Text(text) => self.parse_text_message(&text),
234 Message::Binary(msg) => {
235 let Ok(text) = str::from_utf8(&msg) else {
236 log::warn!(
237 "Received non-UTF-8 BitMEX binary frame ({} bytes)",
238 msg.len()
239 );
240 return None;
241 };
242 self.parse_text_message(text)
243 }
244 Message::Close(_) => {
245 log::debug!("Received close message, waiting for reconnection");
246 None
247 }
248 Message::Ping(data) => {
249 log::trace!("Ping frame with {} bytes (already handled)", data.len());
251 None
252 }
253 Message::Pong(data) => {
254 log::trace!("Received pong frame with {} bytes", data.len());
255 None
256 }
257 Message::Frame(frame) => {
258 log::debug!("Received raw frame: {frame:?}");
259 None
260 }
261 }
262 }
263
264 fn parse_text_message(&self, text: &str) -> Option<BitmexWsFrame> {
265 if text == RECONNECTED {
266 log::info!("Received WebSocket reconnected signal");
267 return Some(BitmexWsFrame::Reconnected);
268 }
269
270 log::trace!("Raw websocket message: {text}");
271
272 if Self::is_heartbeat_message(text) {
273 log::trace!("Ignoring heartbeat control message: {text}");
274 return None;
275 }
276
277 match BitmexTableMessage::from_json_if_table(text) {
278 Ok(Some(table)) => return Some(BitmexWsFrame::Table(table)),
279 Ok(None) => {}
280 Err(e) => {
281 log::error!("Failed to parse WebSocket message: {e}: {text}");
282 return None;
283 }
284 }
285
286 match serde_json::from_str(text) {
287 Ok(msg) => match &msg {
288 BitmexWsFrame::Welcome {
289 version,
290 heartbeat_enabled,
291 limit,
292 ..
293 } => {
294 log::debug!(
295 "Welcome to the BitMEX Realtime API: version={}, heartbeat={}, rate_limit={:?}",
296 version,
297 heartbeat_enabled,
298 limit.as_ref().and_then(|l| l.remaining),
299 );
300 }
301 BitmexWsFrame::Subscription { .. } => return Some(msg),
302 BitmexWsFrame::Error {
303 status,
304 error,
305 request,
306 ..
307 } => {
308 if request
309 .op
310 .eq_ignore_ascii_case(BitmexWsAuthAction::AuthKeyExpires.as_ref())
311 {
312 self.auth_tracker.fail(error.clone());
313 }
314
315 if Self::is_already_subscribed_error(error) {
316 log::debug!(
317 "Ignoring duplicate BitMEX subscription: status={status}, error={error}",
318 );
319 } else {
320 log::error!("Received error from BitMEX: status={status}, error={error}");
321 }
322 }
323 _ => return Some(msg),
324 },
325 Err(e) => {
326 log::error!("Failed to parse WebSocket message: {e}: {text}");
327 }
328 }
329
330 None
331 }
332
333 fn is_heartbeat_message(text: &str) -> bool {
334 let trimmed = text.trim();
335
336 if !trimmed.starts_with('{') || trimmed.len() > 64 {
337 return false;
338 }
339
340 trimmed.contains("\"op\":\"ping\"") || trimmed.contains("\"op\":\"pong\"")
341 }
342
343 fn is_already_subscribed_error(error: &str) -> bool {
344 error.contains("already subscribed to this topic")
345 }
346
347 fn handle_subscription_ack(
348 &self,
349 success: bool,
350 request: Option<&BitmexHttpRequest>,
351 subscribe: Option<&String>,
352 error: Option<&str>,
353 ) {
354 let topics = Self::topics_from_request(request, subscribe);
355
356 if topics.is_empty() {
357 log::debug!("Subscription acknowledgement without topics");
358 return;
359 }
360
361 for topic in topics {
362 if success {
363 self.subscriptions.confirm_subscribe(topic);
364 log::debug!("Subscription confirmed: topic={topic}");
365 } else {
366 self.subscriptions.mark_failure(topic);
367 let reason = error.unwrap_or("Subscription rejected");
368 log::error!("Subscription failed: topic={topic}, error={reason}");
369 }
370 }
371 }
372
373 fn handle_unsubscribe_ack(
374 &self,
375 success: bool,
376 request: Option<&BitmexHttpRequest>,
377 subscribe: Option<&String>,
378 error: Option<&str>,
379 ) {
380 let topics = Self::topics_from_request(request, subscribe);
381
382 if topics.is_empty() {
383 log::debug!("Unsubscription acknowledgement without topics");
384 return;
385 }
386
387 for topic in topics {
388 if success {
389 log::debug!("Unsubscription confirmed: topic={topic}");
390 self.subscriptions.confirm_unsubscribe(topic);
391 } else {
392 let reason = error.unwrap_or("Unsubscription rejected");
393 log::error!(
394 "Unsubscription failed - restoring subscription: topic={topic}, error={reason}",
395 );
396 self.subscriptions.confirm_unsubscribe(topic); self.subscriptions.mark_subscribe(topic); self.subscriptions.confirm_subscribe(topic); }
401 }
402 }
403
404 fn topics_from_request<'a>(
405 request: Option<&'a BitmexHttpRequest>,
406 fallback: Option<&'a String>,
407 ) -> Vec<&'a str> {
408 if let Some(req) = request
409 && !req.args.is_empty()
410 {
411 return req.args.iter().filter_map(|arg| arg.as_str()).collect();
412 }
413
414 fallback.into_iter().map(|topic| topic.as_str()).collect()
415 }
416
417 fn handle_subscription_message(
418 &self,
419 success: bool,
420 subscribe: Option<&String>,
421 request: Option<&BitmexHttpRequest>,
422 error: Option<&str>,
423 ) -> Option<BitmexWsMessage> {
424 if let Some(req) = request {
425 if req
426 .op
427 .eq_ignore_ascii_case(BitmexWsAuthAction::AuthKeyExpires.as_ref())
428 {
429 if success {
430 log::debug!("WebSocket authenticated");
431 self.auth_tracker.succeed();
432 return Some(BitmexWsMessage::Authenticated);
433 } else {
434 let reason = error.unwrap_or("Authentication rejected").to_string();
435 log::error!("WebSocket authentication failed: {reason}");
436 self.auth_tracker.fail(reason);
437 }
438 return None;
439 }
440
441 if req
442 .op
443 .eq_ignore_ascii_case(BitmexWsOperation::Subscribe.as_ref())
444 {
445 self.handle_subscription_ack(success, request, subscribe, error);
446 return None;
447 }
448
449 if req
450 .op
451 .eq_ignore_ascii_case(BitmexWsOperation::Unsubscribe.as_ref())
452 {
453 self.handle_unsubscribe_ack(success, request, subscribe, error);
454 return None;
455 }
456 }
457
458 if subscribe.is_some() {
459 self.handle_subscription_ack(success, request, subscribe, error);
460 return None;
461 }
462
463 if let Some(error) = error {
464 log::warn!("Unhandled subscription control message: success={success}, error={error}");
465 }
466
467 None
468 }
469}
470
471pub(crate) fn should_retry_bitmex_error(error: &BitmexWsError) -> bool {
473 match error {
474 BitmexWsError::TungsteniteError(_) => true, BitmexWsError::ClientError(msg) => {
476 let msg_lower = msg.to_lowercase();
478 msg_lower.contains("timeout")
479 || msg_lower.contains("timed out")
480 || msg_lower.contains("connection")
481 || msg_lower.contains("network")
482 }
483 _ => false,
484 }
485}
486
487pub(crate) fn create_bitmex_timeout_error(msg: String) -> BitmexWsError {
489 BitmexWsError::ClientError(msg)
490}
491
492#[cfg(test)]
493mod tests {
494 use rstest::rstest;
495
496 use super::*;
497 use crate::{
498 common::enums::BitmexOrderStatus,
499 websocket::{
500 enums::BitmexAction,
501 messages::{BitmexTableMessage, OrderData},
502 },
503 };
504
505 fn test_handler() -> BitmexWsFeedHandler {
506 let (_cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
507 let (_raw_tx, raw_rx) = tokio::sync::mpsc::unbounded_channel();
508 let (out_tx, _out_rx) = tokio::sync::mpsc::unbounded_channel();
509
510 BitmexWsFeedHandler::new(
511 Arc::new(AtomicBool::new(false)),
512 cmd_rx,
513 raw_rx,
514 out_tx,
515 AuthTracker::new(),
516 SubscriptionState::new(':'),
517 )
518 }
519
520 #[rstest]
521 #[case(false)]
522 #[case(true)]
523 fn test_json_order_update_routes_from_text_and_binary_frames(#[case] binary: bool) {
524 let json = include_str!("../../test_data/ws_order_update_canceled.json");
525 let message = if binary {
526 Message::Binary(json.as_bytes().to_vec().into())
527 } else {
528 Message::Text(json.into())
529 };
530
531 let handler = test_handler();
532 let Some(BitmexWsFrame::Table(BitmexTableMessage::Order { action, data })) =
533 handler.parse_raw_message(message)
534 else {
535 panic!("expected order table frame");
536 };
537 let OrderData::Update(update) = &data[0] else {
538 panic!("expected sparse order update");
539 };
540
541 assert_eq!(action, BitmexAction::Update);
542 assert_eq!(update.ord_status, Some(BitmexOrderStatus::Canceled));
543 }
544
545 #[rstest]
546 fn test_non_utf8_binary_frame_is_ignored() {
547 let message = Message::Binary(vec![0xFF, 0xFE, 0xFD].into());
548
549 assert!(test_handler().parse_raw_message(message).is_none());
550 }
551
552 #[rstest]
553 fn test_is_heartbeat_message_detection() {
554 assert!(BitmexWsFeedHandler::is_heartbeat_message(
555 "{\"op\":\"ping\"}"
556 ));
557 assert!(BitmexWsFeedHandler::is_heartbeat_message(
558 "{\"op\":\"pong\"}"
559 ));
560 assert!(!BitmexWsFeedHandler::is_heartbeat_message(
561 "{\"op\":\"subscribe\",\"args\":[\"trade:XBTUSD\"]}"
562 ));
563 }
564
565 #[rstest]
566 fn test_is_already_subscribed_error() {
567 let duplicate_error = concat!(
568 "You are already subscribed to this topic:instrument.",
569 " Please see the documentation at https://www.bitmex.com/app/wsAPI."
570 );
571
572 assert!(BitmexWsFeedHandler::is_already_subscribed_error(
573 duplicate_error
574 ));
575 assert!(!BitmexWsFeedHandler::is_already_subscribed_error(
576 "Invalid subscription request"
577 ));
578 }
579}