Skip to main content

nautilus_dydx/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//! Message handler for dYdX WebSocket streams.
17//!
18//! The handler owns the WebSocketClient exclusively and runs in a dedicated
19//! Tokio task within the lock-free I/O boundary. It deserializes raw messages
20//! into venue-specific types without converting to Nautilus domain objects.
21
22use 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/// Commands sent to the feed handler.
57#[derive(Debug, Clone)]
58pub enum HandlerCommand {
59    /// Registers a subscription message for replay.
60    RegisterSubscription {
61        topic: String,
62        subscription: DydxSubscription,
63    },
64    /// Unregisters a subscription message.
65    UnregisterSubscription { topic: String },
66    /// Sends a text message via WebSocket.
67    SendText(String),
68    /// Disconnects the WebSocket client.
69    Disconnect,
70}
71
72/// Deserializes incoming WebSocket messages into venue-specific types.
73///
74/// The handler owns the WebSocketClient exclusively within the lock-free I/O boundary,
75/// eliminating RwLock contention on the hot path.
76pub 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    /// Creates a new [`FeedHandler`].
99    #[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    /// Main processing loop for the handler.
149    ///
150    /// # Panics
151    ///
152    /// This method will not panic. The `expect` call on `iter.next()` is safe
153    /// because we explicitly check that `msgs` is not empty before calling it.
154    pub async fn run(&mut self) {
155        log::debug!("WebSocket handler started");
156
157        loop {
158            // First drain any buffered messages
159            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                        // We just checked that msgs is not empty
181                        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                // Hot path: zero-copy parse for feed messages (orderbook/trades/candles)
219                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                // Cold path: infrequent control messages (connected/subscribed/error)
233                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    /// Handles control messages from the fallback parsing path.
732    ///
733    /// Channel data is handled directly via `handle_feed_message()`.
734    ///
735    /// # Errors
736    ///
737    /// Returns an error if the message cannot be processed.
738    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
784/// Determines if a dYdX WebSocket error should trigger a retry.
785fn 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
805/// Creates a timeout error for the retry manager.
806fn create_dydx_timeout_error(msg: String) -> DydxWsError {
807    DydxWsError::ClientError(msg)
808}