1use std::{
23 collections::VecDeque,
24 fmt::Debug,
25 sync::{
26 Arc,
27 atomic::{AtomicBool, Ordering},
28 },
29};
30
31use ahash::AHashMap;
32use nautilus_network::{
33 RECONNECTED,
34 retry::{RetryManager, create_websocket_retry_manager},
35 websocket::{SubscriptionState, WebSocketClient},
36};
37use tokio_tungstenite::tungstenite::Message;
38use ustr::Ustr;
39
40use super::{
41 DydxWsError, DydxWsResult,
42 client::DYDX_RATE_LIMIT_KEY_SUBSCRIPTION,
43 enums::{DydxWsChannel, DydxWsMessage, DydxWsOutputMessage},
44 error::DydxWebSocketError,
45 messages::{
46 DydxCandle, DydxMarketsContents, DydxOrderbookContents, DydxOrderbookSnapshotContents,
47 DydxSubscription, DydxTradeContents, DydxWsBlockHeightMessage, DydxWsCandlesMessage,
48 DydxWsChannelBatchDataMsg, DydxWsChannelDataMsg, DydxWsConnectedMsg, DydxWsFeedMessage,
49 DydxWsGenericMsg, DydxWsMarketsMessage, DydxWsOrderbookMessage,
50 DydxWsParentSubaccountsMessage, DydxWsSubaccountsChannelContents,
51 DydxWsSubaccountsChannelData, DydxWsSubaccountsMessage, DydxWsSubaccountsSubscribed,
52 DydxWsSubscriptionMsg, DydxWsTradesMessage,
53 },
54};
55
56#[derive(Debug, Clone)]
58pub enum HandlerCommand {
59 RegisterSubscription {
61 topic: String,
62 subscription: DydxSubscription,
63 },
64 UnregisterSubscription { topic: String },
66 SendText(String),
68 Disconnect,
70}
71
72pub struct FeedHandler {
77 cmd_rx: tokio::sync::mpsc::UnboundedReceiver<HandlerCommand>,
78 out_tx: tokio::sync::mpsc::UnboundedSender<DydxWsOutputMessage>,
79 raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
80 client: WebSocketClient,
81 signal: Arc<AtomicBool>,
82 retry_manager: RetryManager<DydxWsError>,
83 subscriptions: SubscriptionState,
84 subscription_messages: AHashMap<String, DydxSubscription>,
85 message_buffer: VecDeque<DydxWsOutputMessage>,
86 book_sequence: AHashMap<String, u64>,
87}
88
89impl Debug for FeedHandler {
90 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
91 f.debug_struct(stringify!(FeedHandler))
92 .field("subscriptions", &self.subscriptions.len())
93 .finish_non_exhaustive()
94 }
95}
96
97impl FeedHandler {
98 #[must_use]
100 pub fn new(
101 cmd_rx: tokio::sync::mpsc::UnboundedReceiver<HandlerCommand>,
102 out_tx: tokio::sync::mpsc::UnboundedSender<DydxWsOutputMessage>,
103 raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
104 client: WebSocketClient,
105 signal: Arc<AtomicBool>,
106 subscriptions: SubscriptionState,
107 ) -> Self {
108 Self {
109 cmd_rx,
110 out_tx,
111 raw_rx,
112 client,
113 signal,
114 retry_manager: create_websocket_retry_manager(),
115 subscriptions,
116 subscription_messages: AHashMap::new(),
117 message_buffer: VecDeque::new(),
118 book_sequence: AHashMap::new(),
119 }
120 }
121
122 async fn send_with_retry(
123 &self,
124 payload: String,
125 rate_limit_keys: Option<&[Ustr]>,
126 ) -> Result<(), DydxWsError> {
127 let keys_owned: Option<Vec<Ustr>> = rate_limit_keys.map(|k| k.to_vec());
128 self.retry_manager
129 .invocation(
130 "websocket_send",
131 || {
132 let payload = payload.clone();
133 let keys = keys_owned.clone();
134 async move {
135 self.client
136 .send_text(payload, keys.as_deref())
137 .await
138 .map_err(|e| DydxWsError::ClientError(format!("Send failed: {e}")))
139 }
140 },
141 should_retry_dydx_error,
142 |e| create_dydx_timeout_error(e.to_string()),
143 )
144 .execute()
145 .await
146 }
147
148 pub async fn run(&mut self) {
155 log::debug!("WebSocket handler started");
156
157 loop {
158 if !self.message_buffer.is_empty() {
160 let msg = self.message_buffer.pop_front().unwrap();
161 if self.out_tx.send(msg).is_err() {
162 log::debug!("Receiver dropped, stopping handler");
163 break;
164 }
165 continue;
166 }
167
168 tokio::select! {
169 Some(cmd) = self.cmd_rx.recv() => {
170 if self.handle_command(cmd).await {
171 break;
172 }
173 }
174
175 Some(msg) = self.raw_rx.recv() => {
176 log::trace!("Handler received raw message");
177 let msgs = self.process_raw_message(msg).await;
178 if !msgs.is_empty() {
179 let mut iter = msgs.into_iter();
180 let first = iter.next().expect("non-empty vec has first element");
182 self.message_buffer.extend(iter);
183 log::trace!("Handler sending message: {:?}", std::mem::discriminant(&first));
184 if self.out_tx.send(first).is_err() {
185 log::debug!("Receiver dropped, stopping handler");
186 break;
187 }
188 }
189 }
190
191 else => {
192 log::debug!("Handler shutting down: channels closed");
193 break;
194 }
195 }
196
197 if self.signal.load(Ordering::Acquire) {
198 log::debug!("Handler received stop signal");
199 break;
200 }
201 }
202 }
203
204 async fn process_raw_message(&mut self, msg: Message) -> Vec<DydxWsOutputMessage> {
205 match msg {
206 Message::Text(txt) => {
207 if txt == RECONNECTED {
208 self.clear_state();
209 self.subscriptions.reset_after_reconnect();
210
211 if let Err(e) = self.replay_subscriptions().await {
212 log::error!("Failed to replay subscriptions after reconnect: {e}");
213 }
214 let topics = self.subscriptions.all_topics();
215 return vec![DydxWsOutputMessage::Reconnected { topics }];
216 }
217
218 match serde_json::from_str::<DydxWsFeedMessage>(&txt) {
220 Ok(feed_msg) => {
221 return self.handle_feed_message(feed_msg);
222 }
223 Err(e) => {
224 if txt.contains("v4_subaccounts") {
225 log::warn!(
226 "[WS_DESER] Failed to parse v4_subaccounts as DydxWsFeedMessage: {e}\nRaw: {txt}"
227 );
228 }
229 }
230 }
231
232 match serde_json::from_str::<serde_json::Value>(&txt) {
234 Ok(val) => match serde_json::from_value::<DydxWsGenericMsg>(val.clone()) {
235 Ok(meta) => {
236 let result = if meta.is_connected() {
237 serde_json::from_value::<DydxWsConnectedMsg>(val)
238 .map(DydxWsMessage::Connected)
239 } else if meta.is_subscribed() {
240 log::debug!("Processing subscribed message via fallback path");
241
242 if let Ok(sub_msg) =
243 serde_json::from_value::<DydxWsSubscriptionMsg>(val.clone())
244 {
245 if sub_msg.channel == DydxWsChannel::Subaccounts {
246 log::debug!("Parsing subaccounts subscription (fallback)");
247 serde_json::from_value::<DydxWsSubaccountsSubscribed>(val)
248 .map(DydxWsMessage::SubaccountsSubscribed)
249 .or_else(|e| {
250 log::warn!(
251 "Failed to parse subaccounts subscription: {e}"
252 );
253 Ok(DydxWsMessage::Subscribed(sub_msg))
254 })
255 } else {
256 Ok(DydxWsMessage::Subscribed(sub_msg))
257 }
258 } else {
259 serde_json::from_value::<DydxWsSubscriptionMsg>(val)
260 .map(DydxWsMessage::Subscribed)
261 }
262 } else if meta.is_unsubscribed() {
263 serde_json::from_value::<DydxWsSubscriptionMsg>(val)
264 .map(DydxWsMessage::Unsubscribed)
265 } else if meta.is_error() {
266 serde_json::from_value::<DydxWebSocketError>(val)
267 .map(DydxWsMessage::Error)
268 } else if meta.is_unknown() {
269 log::warn!("Received unknown WebSocket message type: {txt}",);
270 Ok(DydxWsMessage::Raw(val))
271 } else {
272 Ok(DydxWsMessage::Raw(val))
273 };
274
275 match result {
276 Ok(dydx_msg) => self.handle_dydx_message(dydx_msg).await,
277 Err(e) => {
278 log::error!(
279 "Failed to parse WebSocket message: {e}. Message type: {:?}, Channel: {:?}. Raw: {txt}",
280 meta.msg_type,
281 meta.channel,
282 );
283 vec![]
284 }
285 }
286 }
287 Err(e) => {
288 log::error!(
289 "Failed to parse WebSocket message envelope (DydxWsGenericMsg): {e}\nRaw JSON:\n{txt}"
290 );
291 vec![]
292 }
293 },
294 Err(e) => {
295 let err = DydxWebSocketError::from_message(e.to_string());
296 vec![DydxWsOutputMessage::Error(err)]
297 }
298 }
299 }
300 Message::Pong(_data) => vec![],
301 Message::Ping(_data) => vec![],
302 Message::Binary(_bin) => vec![],
303 Message::Close(_frame) => {
304 log::debug!("WebSocket close frame received");
305 vec![]
306 }
307 Message::Frame(_) => vec![],
308 }
309 }
310
311 async fn handle_dydx_message(&mut self, msg: DydxWsMessage) -> Vec<DydxWsOutputMessage> {
312 match self.handle_message(msg).await {
313 Ok(msgs) => msgs,
314 Err(e) => {
315 log::error!("Error handling message: {e}");
316 vec![]
317 }
318 }
319 }
320
321 fn handle_feed_message(&mut self, feed_msg: DydxWsFeedMessage) -> Vec<DydxWsOutputMessage> {
322 log::trace!(
323 "Handling feed message: {:?}",
324 std::mem::discriminant(&feed_msg)
325 );
326
327 match feed_msg {
328 DydxWsFeedMessage::Subaccounts(msg) => self.handle_subaccounts(msg),
329 DydxWsFeedMessage::Orderbook(msg) => self.handle_orderbook(msg),
330 DydxWsFeedMessage::Trades(msg) => self.handle_trades(msg),
331 DydxWsFeedMessage::Markets(msg) => self.handle_markets_feed(msg),
332 DydxWsFeedMessage::Candles(msg) => self.handle_candles_feed(msg),
333 DydxWsFeedMessage::ParentSubaccounts(msg) => self.handle_parent_subaccounts(msg),
334 DydxWsFeedMessage::BlockHeight(msg) => self.handle_block_height_feed(msg),
335 }
336 }
337
338 fn handle_subaccounts(&self, msg: DydxWsSubaccountsMessage) -> Vec<DydxWsOutputMessage> {
339 match msg {
340 DydxWsSubaccountsMessage::Subscribed(data) => {
341 let topic =
342 self.topic_from_msg(&DydxWsChannel::Subaccounts, &Some(data.id.clone()));
343 self.subscriptions.confirm_subscribe(&topic);
344 log::debug!("Forwarding subaccount subscription to execution client");
345 vec![DydxWsOutputMessage::SubaccountSubscribed(Box::new(data))]
346 }
347 DydxWsSubaccountsMessage::ChannelData(data) => {
348 let has_orders = data.contents.orders.as_ref().is_some_and(|o| !o.is_empty());
349 let has_fills = data.contents.fills.as_ref().is_some_and(|f| !f.is_empty());
350
351 if has_orders || has_fills {
352 log::debug!(
353 "Received {} order(s), {} fill(s) - forwarding to execution client",
354 data.contents.orders.as_ref().map_or(0, |o| o.len()),
355 data.contents.fills.as_ref().map_or(0, |f| f.len())
356 );
357 vec![DydxWsOutputMessage::SubaccountsChannelData(Box::new(data))]
358 } else {
359 vec![]
360 }
361 }
362 DydxWsSubaccountsMessage::Unsubscribed(data) => {
363 let topic = self.topic_from_msg(&DydxWsChannel::Subaccounts, &data.id);
364 self.subscriptions.confirm_unsubscribe(&topic);
365 vec![]
366 }
367 }
368 }
369
370 fn handle_orderbook(&mut self, msg: DydxWsOrderbookMessage) -> Vec<DydxWsOutputMessage> {
371 match msg {
372 DydxWsOrderbookMessage::Subscribed(data) => {
373 let topic = self.topic_from_msg(&DydxWsChannel::Orderbook, &data.id);
374 self.subscriptions.confirm_subscribe(&topic);
375
376 if let Some(id) = &data.id {
377 self.book_sequence.insert(id.clone(), data.message_id);
378 }
379
380 self.deserialize_orderbook_snapshot(&data)
381 }
382 DydxWsOrderbookMessage::ChannelData(data) => {
383 if let Some(id) = &data.id {
384 if let Some(last_id) = self.book_sequence.get(id)
385 && data.message_id <= *last_id
386 {
387 log::warn!(
388 "Orderbook sequence regression for {id}: last {last_id}, received {}",
389 data.message_id
390 );
391 }
392 self.book_sequence.insert(id.clone(), data.message_id);
393 }
394 self.deserialize_orderbook_update(&data)
395 }
396 DydxWsOrderbookMessage::ChannelBatchData(data) => {
397 if let Some(id) = &data.id {
398 if let Some(last_id) = self.book_sequence.get(id)
399 && data.message_id <= *last_id
400 {
401 log::warn!(
402 "Orderbook batch sequence regression for {id}: last {last_id}, received {}",
403 data.message_id
404 );
405 }
406 self.book_sequence.insert(id.clone(), data.message_id);
407 }
408 self.deserialize_orderbook_batch(&data)
409 }
410 DydxWsOrderbookMessage::Unsubscribed(data) => {
411 let topic = self.topic_from_msg(&DydxWsChannel::Orderbook, &data.id);
412 self.subscriptions.confirm_unsubscribe(&topic);
413
414 if let Some(id) = &data.id {
415 self.book_sequence.remove(id);
416 }
417 vec![]
418 }
419 }
420 }
421
422 fn handle_trades(&self, msg: DydxWsTradesMessage) -> Vec<DydxWsOutputMessage> {
423 match msg {
424 DydxWsTradesMessage::Subscribed(data) => {
425 let topic = self.topic_from_msg(&DydxWsChannel::Trades, &data.id);
426 self.subscriptions.confirm_subscribe(&topic);
427 self.deserialize_trades(&data)
428 }
429 DydxWsTradesMessage::ChannelData(data) => self.deserialize_trades(&data),
430 DydxWsTradesMessage::Unsubscribed(data) => {
431 let topic = self.topic_from_msg(&DydxWsChannel::Trades, &data.id);
432 self.subscriptions.confirm_unsubscribe(&topic);
433 vec![]
434 }
435 }
436 }
437
438 fn handle_markets_feed(&self, msg: DydxWsMarketsMessage) -> Vec<DydxWsOutputMessage> {
439 match msg {
440 DydxWsMarketsMessage::Subscribed(data) => {
441 let topic = self.topic_from_msg(&DydxWsChannel::Markets, &data.id);
442 self.subscriptions.confirm_subscribe(&topic);
443 self.deserialize_markets(&data)
444 }
445 DydxWsMarketsMessage::ChannelData(data) => self.deserialize_markets(&data),
446 DydxWsMarketsMessage::Unsubscribed(data) => {
447 let topic = self.topic_from_msg(&DydxWsChannel::Markets, &data.id);
448 self.subscriptions.confirm_unsubscribe(&topic);
449 vec![]
450 }
451 }
452 }
453
454 fn handle_candles_feed(&self, msg: DydxWsCandlesMessage) -> Vec<DydxWsOutputMessage> {
455 match msg {
456 DydxWsCandlesMessage::Subscribed(data) => {
457 let topic = self.topic_from_msg(&DydxWsChannel::Candles, &data.id);
458 self.subscriptions.confirm_subscribe(&topic);
459 vec![]
460 }
461 DydxWsCandlesMessage::ChannelData(data) => self.deserialize_candles(&data),
462 DydxWsCandlesMessage::Unsubscribed(data) => {
463 let topic = self.topic_from_msg(&DydxWsChannel::Candles, &data.id);
464 self.subscriptions.confirm_unsubscribe(&topic);
465 vec![]
466 }
467 }
468 }
469
470 fn handle_parent_subaccounts(
471 &self,
472 msg: DydxWsParentSubaccountsMessage,
473 ) -> Vec<DydxWsOutputMessage> {
474 match msg {
475 DydxWsParentSubaccountsMessage::Subscribed(data) => {
476 let topic = self.topic_from_msg(&DydxWsChannel::ParentSubaccounts, &data.id);
477 self.subscriptions.confirm_subscribe(&topic);
478 self.deserialize_parent_subaccounts(&data)
479 }
480 DydxWsParentSubaccountsMessage::ChannelData(data) => {
481 self.deserialize_parent_subaccounts(&data)
482 }
483 DydxWsParentSubaccountsMessage::Unsubscribed(data) => {
484 let topic = self.topic_from_msg(&DydxWsChannel::ParentSubaccounts, &data.id);
485 self.subscriptions.confirm_unsubscribe(&topic);
486 vec![]
487 }
488 }
489 }
490
491 fn handle_block_height_feed(&self, msg: DydxWsBlockHeightMessage) -> Vec<DydxWsOutputMessage> {
492 match msg {
493 DydxWsBlockHeightMessage::Subscribed(data) => {
494 let topic =
495 self.topic_from_msg(&DydxWsChannel::BlockHeight, &Some(data.id.clone()));
496 self.subscriptions.confirm_subscribe(&topic);
497
498 match data.contents.height.parse::<u64>() {
499 Ok(height) => vec![DydxWsOutputMessage::BlockHeight {
500 height,
501 time: data.contents.time,
502 }],
503 Err(e) => {
504 log::warn!("Failed to parse block height from subscription: {e}");
505 vec![]
506 }
507 }
508 }
509 DydxWsBlockHeightMessage::ChannelData(data) => {
510 match data.contents.block_height.parse::<u64>() {
511 Ok(height) => vec![DydxWsOutputMessage::BlockHeight {
512 height,
513 time: data.contents.time,
514 }],
515 Err(e) => {
516 log::warn!("Failed to parse block height from channel data: {e}");
517 vec![]
518 }
519 }
520 }
521 DydxWsBlockHeightMessage::Unsubscribed(data) => {
522 let topic = self.topic_from_msg(&DydxWsChannel::BlockHeight, &data.id);
523 self.subscriptions.confirm_unsubscribe(&topic);
524 vec![]
525 }
526 }
527 }
528
529 fn deserialize_trades(&self, data: &DydxWsChannelDataMsg) -> Vec<DydxWsOutputMessage> {
530 let Some(id) = data.id.clone() else {
531 log::error!("Missing id for trades channel");
532 return vec![];
533 };
534
535 match serde_json::from_value::<DydxTradeContents>(data.contents.clone()) {
536 Ok(contents) => vec![DydxWsOutputMessage::Trades { id, contents }],
537 Err(e) => {
538 log::error!("Failed to deserialize trade contents: {e}");
539 vec![]
540 }
541 }
542 }
543
544 fn deserialize_orderbook_snapshot(
545 &self,
546 data: &DydxWsChannelDataMsg,
547 ) -> Vec<DydxWsOutputMessage> {
548 let Some(id) = data.id.clone() else {
549 log::error!("Missing id for orderbook snapshot");
550 return vec![];
551 };
552
553 match serde_json::from_value::<DydxOrderbookSnapshotContents>(data.contents.clone()) {
554 Ok(contents) => vec![DydxWsOutputMessage::OrderbookSnapshot { id, contents }],
555 Err(e) => {
556 log::error!("Failed to deserialize orderbook snapshot: {e}");
557 vec![]
558 }
559 }
560 }
561
562 fn deserialize_orderbook_update(
563 &self,
564 data: &DydxWsChannelDataMsg,
565 ) -> Vec<DydxWsOutputMessage> {
566 let Some(id) = data.id.clone() else {
567 log::error!("Missing id for orderbook update");
568 return vec![];
569 };
570
571 match serde_json::from_value::<DydxOrderbookContents>(data.contents.clone()) {
572 Ok(contents) => vec![DydxWsOutputMessage::OrderbookUpdate { id, contents }],
573 Err(e) => {
574 log::error!("Failed to deserialize orderbook contents: {e}");
575 vec![]
576 }
577 }
578 }
579
580 fn deserialize_orderbook_batch(
581 &self,
582 data: &DydxWsChannelBatchDataMsg,
583 ) -> Vec<DydxWsOutputMessage> {
584 let Some(id) = data.id.clone() else {
585 log::error!("Missing id for orderbook batch");
586 return vec![];
587 };
588
589 match serde_json::from_value::<Vec<DydxOrderbookContents>>(data.contents.clone()) {
590 Ok(updates) => vec![DydxWsOutputMessage::OrderbookBatch { id, updates }],
591 Err(e) => {
592 log::error!("Failed to deserialize orderbook batch: {e}");
593 vec![]
594 }
595 }
596 }
597
598 fn deserialize_candles(&self, data: &DydxWsChannelDataMsg) -> Vec<DydxWsOutputMessage> {
599 let Some(id) = data.id.clone() else {
600 log::error!("Missing id for candles channel");
601 return vec![];
602 };
603
604 match serde_json::from_value::<DydxCandle>(data.contents.clone()) {
605 Ok(contents) => vec![DydxWsOutputMessage::Candles { id, contents }],
606 Err(e) => {
607 log::error!("Failed to deserialize candle contents: {e}");
608 vec![]
609 }
610 }
611 }
612
613 fn deserialize_markets(&self, data: &DydxWsChannelDataMsg) -> Vec<DydxWsOutputMessage> {
614 match serde_json::from_value::<DydxMarketsContents>(data.contents.clone()) {
615 Ok(contents) => vec![DydxWsOutputMessage::Markets(contents)],
616 Err(e) => {
617 log::error!("Failed to deserialize markets contents: {e}");
618 vec![]
619 }
620 }
621 }
622
623 fn deserialize_parent_subaccounts(
624 &self,
625 data: &DydxWsChannelDataMsg,
626 ) -> Vec<DydxWsOutputMessage> {
627 match serde_json::from_value::<DydxWsSubaccountsChannelContents>(data.contents.clone()) {
628 Ok(contents) => {
629 let has_orders = contents.orders.as_ref().is_some_and(|o| !o.is_empty());
630 let has_fills = contents.fills.as_ref().is_some_and(|f| !f.is_empty());
631
632 if has_orders || has_fills {
633 let channel_data = DydxWsSubaccountsChannelData {
634 connection_id: data.connection_id.clone(),
635 message_id: data.message_id,
636 id: data.id.clone().unwrap_or_default(),
637 version: data.version.clone().unwrap_or_default(),
638 contents,
639 };
640 vec![DydxWsOutputMessage::SubaccountsChannelData(Box::new(
641 channel_data,
642 ))]
643 } else {
644 vec![]
645 }
646 }
647 Err(e) => {
648 log::error!("Failed to deserialize parent subaccounts contents: {e}");
649 vec![]
650 }
651 }
652 }
653
654 async fn handle_command(&mut self, command: HandlerCommand) -> bool {
655 match command {
656 HandlerCommand::RegisterSubscription {
657 topic,
658 subscription,
659 } => {
660 self.subscription_messages.insert(topic, subscription);
661 }
662 HandlerCommand::UnregisterSubscription { topic } => {
663 self.subscription_messages.remove(&topic);
664 }
665 HandlerCommand::SendText(text) => {
666 if let Err(e) = self
667 .send_with_retry(text, Some(DYDX_RATE_LIMIT_KEY_SUBSCRIPTION.as_slice()))
668 .await
669 {
670 log::error!("Failed to send WebSocket text after retries: {e}");
671 }
672 }
673 HandlerCommand::Disconnect => {
674 log::debug!("Disconnect command received");
675 self.client.disconnect().await;
676 return true;
677 }
678 }
679 false
680 }
681
682 fn topic_from_msg(&self, channel: &DydxWsChannel, id: &Option<String>) -> String {
683 if matches!(channel, DydxWsChannel::BlockHeight) {
684 return channel.as_ref().to_string();
685 }
686
687 match id {
688 Some(id) => format!(
689 "{}{}{}",
690 channel.as_ref(),
691 self.subscriptions.delimiter(),
692 id
693 ),
694 None => channel.as_ref().to_string(),
695 }
696 }
697
698 fn clear_state(&mut self) {
699 let buffer_count = self.message_buffer.len();
700 let seq_count = self.book_sequence.len();
701 self.message_buffer.clear();
702 self.book_sequence.clear();
703 log::debug!(
704 "Cleared reconnect state: message_buffer={buffer_count}, book_sequence={seq_count}"
705 );
706 }
707
708 async fn replay_subscriptions(&self) -> DydxWsResult<()> {
709 let topics = self.subscriptions.all_topics();
710 for topic in topics {
711 let Some(subscription) = self.subscription_messages.get(&topic).cloned() else {
712 log::warn!("No preserved subscription message for topic: {topic}");
713 continue;
714 };
715
716 let payload = serde_json::to_string(&subscription)?;
717 self.subscriptions.mark_subscribe(&topic);
718
719 if let Err(e) = self
720 .send_with_retry(payload, Some(DYDX_RATE_LIMIT_KEY_SUBSCRIPTION.as_slice()))
721 .await
722 {
723 self.subscriptions.mark_failure(&topic);
724 return Err(e);
725 }
726 }
727
728 Ok(())
729 }
730
731 pub async fn handle_message(
739 &mut self,
740 msg: DydxWsMessage,
741 ) -> DydxWsResult<Vec<DydxWsOutputMessage>> {
742 match msg {
743 DydxWsMessage::Connected(_) => {
744 log::debug!("dYdX WebSocket connected");
745 Ok(vec![])
746 }
747 DydxWsMessage::Subscribed(sub) => {
748 log::debug!("Subscribed to {} (id: {:?})", sub.channel, sub.id);
749 let topic = self.topic_from_msg(&sub.channel, &sub.id);
750 self.subscriptions.confirm_subscribe(&topic);
751 Ok(vec![])
752 }
753 DydxWsMessage::SubaccountsSubscribed(msg) => {
754 log::debug!("Subaccounts subscribed with initial state (fallback path)");
755 let topic = self.topic_from_msg(&DydxWsChannel::Subaccounts, &Some(msg.id.clone()));
756 self.subscriptions.confirm_subscribe(&topic);
757 Ok(vec![DydxWsOutputMessage::SubaccountSubscribed(Box::new(
758 msg,
759 ))])
760 }
761 DydxWsMessage::Unsubscribed(unsub) => {
762 log::debug!("Unsubscribed from {} (id: {:?})", unsub.channel, unsub.id);
763 let topic = self.topic_from_msg(&unsub.channel, &unsub.id);
764 self.subscriptions.confirm_unsubscribe(&topic);
765 Ok(vec![])
766 }
767 DydxWsMessage::Error(err) => Ok(vec![DydxWsOutputMessage::Error(err)]),
768 DydxWsMessage::Reconnected => {
769 self.clear_state();
770 self.subscriptions.reset_after_reconnect();
771
772 if let Err(e) = self.replay_subscriptions().await {
773 log::error!("Failed to replay subscriptions after reconnect message: {e}");
774 }
775 let topics = self.subscriptions.all_topics();
776 Ok(vec![DydxWsOutputMessage::Reconnected { topics }])
777 }
778 DydxWsMessage::Pong => Ok(vec![]),
779 DydxWsMessage::Raw(_) => Ok(vec![]),
780 }
781 }
782}
783
784fn should_retry_dydx_error(error: &DydxWsError) -> bool {
786 match error {
787 DydxWsError::Transport(_) => true,
788 DydxWsError::Send(_) => true,
789 DydxWsError::ClientError(msg) => {
790 let msg_lower = msg.to_lowercase();
791 msg_lower.contains("timeout")
792 || msg_lower.contains("timed out")
793 || msg_lower.contains("connection")
794 || msg_lower.contains("network")
795 }
796 DydxWsError::NotConnected
797 | DydxWsError::Json(_)
798 | DydxWsError::Parse(_)
799 | DydxWsError::Authentication(_)
800 | DydxWsError::Subscription(_)
801 | DydxWsError::Venue(_) => false,
802 }
803}
804
805fn create_dydx_timeout_error(msg: String) -> DydxWsError {
807 DydxWsError::ClientError(msg)
808}