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