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