nautilus_bitmex/websocket/
handler.rs1use 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#[derive(Debug)]
39pub enum HandlerCommand {
40 SetClient(WebSocketClient),
42 Disconnect,
44 Authenticate { payload: SecretString },
46 Subscribe { topics: Vec<String> },
48 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 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 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 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 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 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 self.subscriptions.confirm_unsubscribe(topic); self.subscriptions.mark_subscribe(topic); self.subscriptions.confirm_subscribe(topic); }
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
480pub(crate) fn should_retry_bitmex_error(error: &BitmexWsError) -> bool {
482 match error {
483 BitmexWsError::TungsteniteError(_) => true, BitmexWsError::ClientError(msg) => {
485 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
496pub(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}