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            .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    /// Main processing loop for the handler.
148    ///
149    /// # Panics
150    ///
151    /// This method will not panic. The `expect` call on `iter.next()` is safe
152    /// because we explicitly check that `msgs` is not empty before calling it.
153    pub async fn run(&mut self) {
154        log::debug!("WebSocket handler started");
155
156        loop {
157            // First drain any buffered messages
158            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                        // We just checked that msgs is not empty
180                        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                // Hot path: zero-copy parse for feed messages (orderbook/trades/candles)
218                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                // Cold path: infrequent control messages (connected/subscribed/error)
232                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    /// Handles control messages from the fallback parsing path.
731    ///
732    /// Channel data is handled directly via `handle_feed_message()`.
733    ///
734    /// # Errors
735    ///
736    /// Returns an error if the message cannot be processed.
737    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
783/// Determines if a dYdX WebSocket error should trigger a retry.
784fn 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
804/// Creates a timeout error for the retry manager.
805fn create_dydx_timeout_error(msg: String) -> DydxWsError {
806    DydxWsError::ClientError(msg)
807}