Skip to main content

nautilus_bitmex/websocket/
handler.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! WebSocket message handler for BitMEX.
17
18use 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/// Commands sent from the outer client to the inner message handler.
37#[derive(Debug)]
38pub enum HandlerCommand {
39    /// Set the WebSocketClient for the handler to use.
40    SetClient(WebSocketClient),
41    /// Disconnect the WebSocket connection.
42    Disconnect,
43    /// Send authentication payload to the WebSocket.
44    Authenticate { payload: String },
45    /// Subscribe to the given topics.
46    Subscribe { topics: Vec<String> },
47    /// Unsubscribe from the given topics.
48    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    /// Creates a new [`BitmexWsFeedHandler`] instance.
64    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    /// Sends a WebSocket message with retry logic.
93    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                    // Handle ping frames directly for minimal latency
175                    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                // Handle shutdown - either channel closed or stream ended
223                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                // Handled in select! loop before parse_raw_message
250                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                // Venue rejected unsubscribe, so we're still subscribed. Restore state:
397                self.subscriptions.confirm_unsubscribe(topic); // Clear pending_unsubscribe
398                self.subscriptions.mark_subscribe(topic); // Mark as subscribing
399                self.subscriptions.confirm_subscribe(topic); // Confirm subscription
400            }
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
471/// Returns `true` when a BitMEX error should be retried.
472pub(crate) fn should_retry_bitmex_error(error: &BitmexWsError) -> bool {
473    match error {
474        BitmexWsError::TungsteniteError(_) => true, // Network errors are retryable
475        BitmexWsError::ClientError(msg) => {
476            // Retry on timeout and connection errors (case-insensitive)
477            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
487/// Creates a timeout error for BitMEX retry logic.
488pub(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}