Skip to main content

nautilus_dydx/
data.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//! Live market data client implementation for the dYdX adapter.
17
18use std::{
19    str::FromStr,
20    sync::{
21        Arc,
22        atomic::{AtomicBool, Ordering},
23    },
24    time::Duration,
25};
26
27use anyhow::Context;
28use dashmap::DashMap;
29use futures_util::{Stream, StreamExt, pin_mut};
30use nautilus_common::{
31    clients::DataClient,
32    live::{runner::get_data_event_sender, sender::EventSender},
33    messages::{
34        DataEvent, DataResponse,
35        data::{
36            BarsResponse, BookResponse, FundingRatesResponse, InstrumentResponse,
37            InstrumentsResponse, RequestBars, RequestBookSnapshot, RequestFundingRates,
38            RequestInstrument, RequestInstruments, RequestTrades, SubscribeBars,
39            SubscribeBookDeltas, SubscribeFundingRates, SubscribeIndexPrices, SubscribeInstrument,
40            SubscribeInstrumentStatus, SubscribeInstruments, SubscribeMarkPrices, SubscribeQuotes,
41            SubscribeTrades, TradesResponse, UnsubscribeBars, UnsubscribeBookDeltas,
42            UnsubscribeFundingRates, UnsubscribeIndexPrices, UnsubscribeInstrument,
43            UnsubscribeInstrumentStatus, UnsubscribeInstruments, UnsubscribeMarkPrices,
44            UnsubscribeQuotes, UnsubscribeTrades,
45        },
46    },
47};
48use nautilus_core::{
49    AtomicMap, AtomicSet,
50    datetime::datetime_to_unix_nanos,
51    time::{AtomicTime, get_atomic_clock_realtime},
52};
53use nautilus_live::{
54    SocketControlFactory,
55    task::{TaskGroup, TaskGroupGuard, TaskSpawner},
56};
57use nautilus_model::{
58    data::{
59        Bar, BarSpecification, BarType, BookOrder, Data as NautilusData, FundingRateUpdate,
60        IndexPriceUpdate, InstrumentStatus, MarkPriceUpdate, OrderBookDelta, OrderBookDeltas,
61        QuoteTick,
62    },
63    enums::{BookAction, BookType, MarketStatusAction, OrderSide, RecordFlag},
64    identifiers::{ClientId, InstrumentId, Symbol, Venue},
65    instruments::{Instrument, InstrumentAny},
66    orderbook::OrderBook,
67    types::Quantity,
68};
69use rust_decimal::Decimal;
70use ustr::Ustr;
71
72use crate::{
73    common::{
74        consts::DYDX_VENUE,
75        enums::DydxCandleResolution,
76        instrument_cache::InstrumentCache,
77        parse::{extract_raw_symbol, parse_price},
78    },
79    config::DydxDataClientConfig,
80    http::client::DydxHttpClient,
81    websocket::{
82        client::{DydxWebSocketClient, candle_ids_from_topics},
83        enums::DydxWsOutputMessage,
84        parse as ws_parse,
85    },
86};
87
88/// dYdX data client for live market data streaming and historical data requests.
89///
90/// This client integrates with the Nautilus DataEngine to provide:
91/// - Real-time market data via WebSocket subscriptions
92/// - Historical data via REST API requests
93/// - Automatic instrument discovery and caching
94/// - Connection lifecycle management
95#[derive(Debug)]
96pub struct DydxDataClient {
97    clock: &'static AtomicTime,
98    client_id: ClientId,
99    config: DydxDataClientConfig,
100    http_client: DydxHttpClient,
101    ws_client: DydxWebSocketClient,
102    is_connected: AtomicBool,
103    session_tasks: TaskGroup,
104    command_tasks: TaskGroup,
105    shutdown_errors: Vec<String>,
106    data_sender: EventSender<DataEvent>,
107    instrument_cache: Arc<InstrumentCache>,
108    order_books: Arc<DashMap<InstrumentId, OrderBook>>,
109    last_quotes: Arc<DashMap<InstrumentId, QuoteTick>>,
110    incomplete_bars: Arc<DashMap<BarType, Bar>>,
111    bar_type_mappings: Arc<AtomicMap<String, BarType>>,
112    active_quote_subs: Arc<AtomicSet<InstrumentId>>,
113    active_delta_subs: Arc<AtomicSet<InstrumentId>>,
114    active_trade_subs: Arc<AtomicSet<InstrumentId>>,
115    active_bar_subs: Arc<AtomicMap<(InstrumentId, String), BarType>>,
116    active_mark_price_subs: Arc<AtomicSet<InstrumentId>>,
117    active_index_price_subs: Arc<AtomicSet<InstrumentId>>,
118    active_funding_rate_subs: Arc<AtomicSet<InstrumentId>>,
119    active_instrument_status_subs: Arc<AtomicSet<InstrumentId>>,
120    last_instrument_statuses: Arc<DashMap<InstrumentId, InstrumentStatus>>,
121}
122
123impl DydxDataClient {
124    fn map_bar_spec_to_resolution(spec: &BarSpecification) -> anyhow::Result<&'static str> {
125        let resolution: &'static str = DydxCandleResolution::from_bar_spec(spec)?.into();
126        Ok(resolution)
127    }
128
129    /// Creates a new [`DydxDataClient`] instance.
130    ///
131    /// # Errors
132    ///
133    /// Returns an error if the client fails to initialize.
134    pub fn new(
135        client_id: ClientId,
136        config: DydxDataClientConfig,
137        http_client: DydxHttpClient,
138        ws_client: DydxWebSocketClient,
139    ) -> anyhow::Result<Self> {
140        let clock = get_atomic_clock_realtime();
141        let data_sender = get_data_event_sender();
142        let ws_client =
143            ws_client.with_socket_factory(SocketControlFactory::new(client_id, Some(*DYDX_VENUE)));
144
145        let instrument_cache = Arc::clone(http_client.instrument_cache());
146        let session_tasks = TaskGroup::new();
147        let command_tasks = TaskGroup::new();
148
149        Ok(Self {
150            clock,
151            client_id,
152            config,
153            http_client,
154            ws_client,
155            is_connected: AtomicBool::new(false),
156            session_tasks,
157            command_tasks,
158            shutdown_errors: Vec::new(),
159            data_sender,
160            instrument_cache,
161            order_books: Arc::new(DashMap::new()),
162            last_quotes: Arc::new(DashMap::new()),
163            incomplete_bars: Arc::new(DashMap::new()),
164            bar_type_mappings: Arc::new(AtomicMap::new()),
165            active_quote_subs: Arc::new(AtomicSet::new()),
166            active_delta_subs: Arc::new(AtomicSet::new()),
167            active_trade_subs: Arc::new(AtomicSet::new()),
168            active_bar_subs: Arc::new(AtomicMap::new()),
169            active_mark_price_subs: Arc::new(AtomicSet::new()),
170            active_index_price_subs: Arc::new(AtomicSet::new()),
171            active_funding_rate_subs: Arc::new(AtomicSet::new()),
172            active_instrument_status_subs: Arc::new(AtomicSet::new()),
173            last_instrument_statuses: Arc::new(DashMap::new()),
174        })
175    }
176
177    /// Returns the venue for this data client.
178    #[must_use]
179    pub fn venue(&self) -> Venue {
180        *DYDX_VENUE
181    }
182
183    /// Returns a reference to the client configuration.
184    #[must_use]
185    pub fn config(&self) -> &DydxDataClientConfig {
186        &self.config
187    }
188
189    /// Returns `true` when the client is connected.
190    #[must_use]
191    pub fn is_connected(&self) -> bool {
192        self.is_connected.load(Ordering::Relaxed)
193    }
194
195    fn spawn_ws<F>(&self, fut: F, context: &'static str)
196    where
197        F: std::future::Future<Output = anyhow::Result<()>> + Send + 'static,
198    {
199        let future = async move {
200            if let Err(e) = fut.await {
201                log::error!("{context}: {e:?}");
202            }
203        };
204
205        if let Err(e) = self.command_tasks.spawn(future) {
206            log::warn!("Skipping dYdX {context} after shutdown began: {e}");
207        }
208    }
209
210    fn spawn_command<F>(&self, future: F)
211    where
212        F: std::future::Future<Output = ()> + Send + 'static,
213    {
214        if let Err(e) = self.command_tasks.spawn(future) {
215            log::warn!("Skipping dYdX data command after shutdown began: {e}");
216        }
217    }
218
219    fn spawn_ws_stream_handler(
220        &self,
221        stream: impl Stream<Item = DydxWsOutputMessage> + Send + 'static,
222        ctx: WsMessageContext,
223    ) -> anyhow::Result<()> {
224        let cancellation = self.session_tasks.cancellation_token();
225
226        let future = async move {
227            log::debug!("Message processing task started");
228            pin_mut!(stream);
229
230            loop {
231                tokio::select! {
232                    maybe_msg = stream.next() => {
233                        match maybe_msg {
234                            Some(msg) => Self::handle_ws_message(msg, &ctx),
235                            None => {
236                                log::debug!("WebSocket message channel closed");
237                                break;
238                            }
239                        }
240                    }
241                    () = cancellation.cancelled() => {
242                        log::debug!("WebSocket message task cancelled");
243                        break;
244                    }
245                }
246            }
247            log::debug!("WebSocket stream handler ended");
248        };
249
250        self.session_tasks
251            .spawn(future)
252            .context("failed to register dYdX WebSocket stream task")?;
253        Ok(())
254    }
255
256    async fn finish_tasks(&self) -> anyhow::Result<()> {
257        let (session_result, command_result) = tokio::join!(
258            self.session_tasks
259                .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
260            self.command_tasks
261                .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
262        );
263        session_result.context("failed to finish dYdX data session tasks")?;
264        command_result.context("failed to finish dYdX data command tasks")?;
265        Ok(())
266    }
267
268    async fn prepare_task_groups(&mut self) -> anyhow::Result<()> {
269        if !self.session_tasks.is_open() || !self.command_tasks.is_open() {
270            self.session_tasks.begin_shutdown();
271            self.command_tasks.begin_shutdown();
272            self.ws_client.begin_shutdown();
273            self.finish_shutdown().await?;
274            self.session_tasks
275                .start_generation()
276                .context("failed to start dYdX data session task generation")?;
277            self.command_tasks
278                .start_generation()
279                .context("failed to start dYdX data command task generation")?;
280        }
281        Ok(())
282    }
283
284    async fn finish_shutdown(&mut self) -> anyhow::Result<()> {
285        if let Err(e) = self
286            .ws_client
287            .disconnect()
288            .await
289            .context("failed to disconnect dYdX websocket")
290        {
291            self.shutdown_errors.push(e.to_string());
292        }
293
294        if let Err(e) = self.finish_tasks().await {
295            self.shutdown_errors.push(e.to_string());
296        }
297
298        if !self.shutdown_errors.is_empty() {
299            anyhow::bail!(std::mem::take(&mut self.shutdown_errors).join("; "));
300        }
301        Ok(())
302    }
303
304    async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
305        self.session_tasks.begin_shutdown();
306        self.command_tasks.begin_shutdown();
307        self.ws_client.begin_shutdown();
308        let shutdown_result = self.finish_shutdown().await;
309        self.is_connected.store(false, Ordering::Release);
310        shutdown_result
311    }
312
313    async fn bootstrap_instruments(&self) -> anyhow::Result<Vec<InstrumentAny>> {
314        self.http_client
315            .fetch_and_cache_instruments()
316            .await
317            .context("failed to load instruments from dYdX")?;
318
319        let instruments: Vec<InstrumentAny> = self.http_client.all_instruments();
320
321        if instruments.is_empty() {
322            log::warn!("No instruments were loaded");
323            return Ok(instruments);
324        }
325
326        log::debug!("Loaded {} instruments into shared cache", instruments.len());
327
328        self.ws_client.cache_instruments(instruments.clone());
329
330        for instrument in &instruments {
331            if let Err(e) = self
332                .data_sender
333                .send(DataEvent::Instrument(instrument.clone()))
334            {
335                log::warn!("Failed to publish instrument {}: {e}", instrument.id());
336            }
337        }
338        log::debug!("Published {} instruments to data engine", instruments.len());
339
340        Ok(instruments)
341    }
342}
343
344#[async_trait::async_trait(?Send)]
345impl DataClient for DydxDataClient {
346    fn client_id(&self) -> ClientId {
347        self.client_id
348    }
349
350    fn venue(&self) -> Option<Venue> {
351        Some(*DYDX_VENUE)
352    }
353
354    fn start(&mut self) -> anyhow::Result<()> {
355        log::info!(
356            "Starting: client_id={}, is_testnet={}",
357            self.client_id,
358            self.http_client.is_testnet()
359        );
360        Ok(())
361    }
362
363    fn stop(&mut self) -> anyhow::Result<()> {
364        log::info!("Stopping {}", self.client_id);
365        self.session_tasks.begin_shutdown();
366        self.command_tasks.begin_shutdown();
367        self.ws_client.begin_shutdown();
368        self.is_connected.store(false, Ordering::Relaxed);
369        Ok(())
370    }
371
372    fn reset(&mut self) -> anyhow::Result<()> {
373        log::debug!("Resetting {}", self.client_id);
374        self.session_tasks.begin_shutdown();
375        self.command_tasks.begin_shutdown();
376        self.ws_client.begin_shutdown();
377        self.is_connected.store(false, Ordering::Relaxed);
378        Ok(())
379    }
380
381    fn dispose(&mut self) -> anyhow::Result<()> {
382        log::debug!("Disposing {}", self.client_id);
383        self.stop()
384    }
385
386    async fn connect(&mut self) -> anyhow::Result<()> {
387        if self.is_connected() && self.session_tasks.is_open() && self.command_tasks.is_open() {
388            return Ok(());
389        }
390
391        log::info!("Connecting");
392
393        self.prepare_task_groups().await?;
394        let ws_client = self.ws_client.clone();
395        let setup_guard =
396            TaskGroupGuard::new(&[&self.session_tasks, &self.command_tasks], move || {
397                ws_client.begin_shutdown();
398            });
399
400        self.bootstrap_instruments().await?;
401
402        let session_result = async {
403            self.ws_client
404                .connect()
405                .await
406                .context("failed to connect dYdX websocket")?;
407
408            self.ws_client
409                .subscribe_markets()
410                .await
411                .context("failed to subscribe to markets channel")?;
412
413            let seen_tickers: Arc<AtomicSet<Ustr>> = Arc::new(AtomicSet::new());
414
415            for instrument in self.instrument_cache.all_instruments() {
416                let id = instrument.id();
417                let ticker = extract_raw_symbol(id.symbol.as_str());
418                seen_tickers.insert(Ustr::from(ticker));
419            }
420
421            let command_spawner = self
422                .command_tasks
423                .spawner()
424                .context("dYdX data command task admission is closed")?;
425            let ctx = WsMessageContext {
426                clock: self.clock,
427                data_sender: self.data_sender.clone(),
428                instrument_cache: self.instrument_cache.clone(),
429                order_books: self.order_books.clone(),
430                last_quotes: self.last_quotes.clone(),
431                ws_client: self.ws_client.clone(),
432                http_client: self.http_client.clone(),
433                active_quote_subs: self.active_quote_subs.clone(),
434                active_delta_subs: self.active_delta_subs.clone(),
435                active_trade_subs: self.active_trade_subs.clone(),
436                active_bar_subs: self.active_bar_subs.clone(),
437                incomplete_bars: self.incomplete_bars.clone(),
438                bar_type_mappings: self.bar_type_mappings.clone(),
439                active_mark_price_subs: self.active_mark_price_subs.clone(),
440                active_index_price_subs: self.active_index_price_subs.clone(),
441                active_funding_rate_subs: self.active_funding_rate_subs.clone(),
442                active_instrument_status_subs: self.active_instrument_status_subs.clone(),
443                last_instrument_statuses: self.last_instrument_statuses.clone(),
444                bars_timestamp_on_close: self.ws_client.bars_timestamp_on_close(),
445                pending_bars: Arc::new(DashMap::new()),
446                seen_tickers,
447                command_spawner,
448            };
449
450            let stream = self.ws_client.stream();
451            self.spawn_ws_stream_handler(stream, ctx)?;
452
453            Ok::<(), anyhow::Error>(())
454        }
455        .await;
456
457        if let Err(e) = session_result {
458            if let Err(teardown_error) = self.teardown_partial_connect().await {
459                return Err(e.context(format!(
460                    "dYdX data startup teardown failed: {teardown_error}"
461                )));
462            }
463            return Err(e);
464        }
465
466        self.is_connected.store(true, Ordering::Relaxed);
467        setup_guard.disarm();
468        log::info!("Connected");
469
470        Ok(())
471    }
472
473    async fn disconnect(&mut self) -> anyhow::Result<()> {
474        log::info!("Disconnecting");
475
476        self.teardown_partial_connect().await?;
477
478        self.last_instrument_statuses.clear();
479        self.is_connected.store(false, Ordering::Relaxed);
480        log::info!("Disconnected dYdX data client");
481
482        Ok(())
483    }
484
485    fn is_connected(&self) -> bool {
486        self.is_connected.load(Ordering::Relaxed)
487    }
488
489    fn is_disconnected(&self) -> bool {
490        !self.is_connected()
491    }
492
493    fn subscribe_instruments(&mut self, _cmd: SubscribeInstruments) -> anyhow::Result<()> {
494        log::debug!(
495            "subscribe_instruments: dYdX instruments discovered via global v4_markets channel"
496        );
497        Ok(())
498    }
499
500    fn subscribe_instrument(&mut self, cmd: SubscribeInstrument) -> anyhow::Result<()> {
501        if let Some(instrument) = self.instrument_cache.get(&cmd.instrument_id) {
502            log::debug!("Sending cached instrument for {}", cmd.instrument_id);
503            if let Err(e) = self.data_sender.send(DataEvent::Instrument(instrument)) {
504                log::warn!("Failed to send instrument {}: {e}", cmd.instrument_id);
505            }
506        } else {
507            log::warn!(
508                "Instrument {} not found in cache (available: {})",
509                cmd.instrument_id,
510                self.instrument_cache.len()
511            );
512        }
513        Ok(())
514    }
515
516    fn subscribe_book_deltas(&mut self, cmd: SubscribeBookDeltas) -> anyhow::Result<()> {
517        if cmd.book_type != BookType::L2_MBP {
518            anyhow::bail!(
519                "dYdX only supports L2_MBP order book deltas, received {:?}",
520                cmd.book_type
521            );
522        }
523
524        self.ensure_order_book(cmd.instrument_id, BookType::L2_MBP);
525        self.active_delta_subs.insert(cmd.instrument_id);
526
527        let ws = self.ws_client.clone();
528        let instrument_id = cmd.instrument_id;
529
530        self.spawn_ws(
531            async move {
532                ws.subscribe_orderbook(instrument_id)
533                    .await
534                    .context("orderbook subscription")
535            },
536            "dYdX orderbook subscription",
537        );
538
539        Ok(())
540    }
541
542    fn subscribe_quotes(&mut self, cmd: SubscribeQuotes) -> anyhow::Result<()> {
543        log::debug!(
544            "Subscribe_quotes for {}: subscribing to orderbook WS channel for quote synthesis",
545            cmd.instrument_id
546        );
547
548        self.ensure_order_book(cmd.instrument_id, BookType::L2_MBP);
549        self.active_quote_subs.insert(cmd.instrument_id);
550        let ws = self.ws_client.clone();
551        let instrument_id = cmd.instrument_id;
552
553        self.spawn_ws(
554            async move {
555                ws.subscribe_orderbook(instrument_id)
556                    .await
557                    .context("orderbook subscription (for quotes)")
558            },
559            "dYdX orderbook subscription (quotes)",
560        );
561
562        Ok(())
563    }
564
565    fn subscribe_trades(&mut self, cmd: SubscribeTrades) -> anyhow::Result<()> {
566        let ws = self.ws_client.clone();
567        let instrument_id = cmd.instrument_id;
568
569        self.active_trade_subs.insert(instrument_id);
570
571        self.spawn_ws(
572            async move {
573                ws.subscribe_trades(instrument_id)
574                    .await
575                    .context("trade subscription")
576            },
577            "dYdX trade subscription",
578        );
579
580        Ok(())
581    }
582
583    fn subscribe_mark_prices(&mut self, cmd: SubscribeMarkPrices) -> anyhow::Result<()> {
584        let instrument_id = cmd.instrument_id;
585        self.active_mark_price_subs.insert(instrument_id);
586        log::debug!("Subscribed to mark prices for {instrument_id} (via v4_markets channel)");
587        Ok(())
588    }
589
590    fn subscribe_index_prices(&mut self, cmd: SubscribeIndexPrices) -> anyhow::Result<()> {
591        let instrument_id = cmd.instrument_id;
592        self.active_index_price_subs.insert(instrument_id);
593        log::debug!("Subscribed to index prices for {instrument_id} (via v4_markets channel)");
594        Ok(())
595    }
596
597    fn subscribe_bars(&mut self, cmd: SubscribeBars) -> anyhow::Result<()> {
598        let ws = self.ws_client.clone();
599        let instrument_id = cmd.bar_type.instrument_id();
600        let spec = cmd.bar_type.spec();
601
602        let resolution = Self::map_bar_spec_to_resolution(&spec)?;
603        let bar_type = cmd.bar_type;
604        self.active_bar_subs
605            .insert((instrument_id, resolution.to_string()), bar_type);
606
607        let ticker = extract_raw_symbol(instrument_id.symbol.as_str());
608        let topic = format!("{ticker}/{resolution}");
609        self.bar_type_mappings.insert(topic, bar_type);
610
611        self.spawn_ws(
612            async move {
613                ws.subscribe_candles(instrument_id, resolution)
614                    .await
615                    .context("candles subscription")
616            },
617            "dYdX candles subscription",
618        );
619
620        Ok(())
621    }
622
623    fn subscribe_funding_rates(&mut self, cmd: SubscribeFundingRates) -> anyhow::Result<()> {
624        let instrument_id = cmd.instrument_id;
625        self.active_funding_rate_subs.insert(instrument_id);
626        log::debug!("Subscribed to funding rates for {instrument_id} (via v4_markets channel)");
627        Ok(())
628    }
629
630    fn subscribe_instrument_status(
631        &mut self,
632        cmd: SubscribeInstrumentStatus,
633    ) -> anyhow::Result<()> {
634        let instrument_id = cmd.instrument_id;
635        self.active_instrument_status_subs.insert(instrument_id);
636        log::debug!("Subscribed to instrument status for {instrument_id} (via v4_markets channel)");
637
638        // Replay last known status (initial snapshot arrives before subscription)
639        if let Some(status) = self.last_instrument_statuses.get(&instrument_id)
640            && let Err(e) = self.data_sender.send(DataEvent::InstrumentStatus(*status))
641        {
642            log::error!("Failed to replay instrument status for {instrument_id}: {e}");
643        }
644
645        Ok(())
646    }
647
648    fn unsubscribe_instruments(&mut self, _cmd: &UnsubscribeInstruments) -> anyhow::Result<()> {
649        log::debug!("unsubscribe_instruments: dYdX markets channel is global; no-op");
650        Ok(())
651    }
652
653    fn unsubscribe_instrument(&mut self, _cmd: &UnsubscribeInstrument) -> anyhow::Result<()> {
654        log::debug!("unsubscribe_instrument: dYdX markets channel is global; no-op");
655        Ok(())
656    }
657
658    fn unsubscribe_book_deltas(&mut self, cmd: &UnsubscribeBookDeltas) -> anyhow::Result<()> {
659        self.active_delta_subs.remove(&cmd.instrument_id);
660
661        let ws = self.ws_client.clone();
662        let instrument_id = cmd.instrument_id;
663
664        self.spawn_ws(
665            async move {
666                ws.unsubscribe_orderbook(instrument_id)
667                    .await
668                    .context("orderbook unsubscription")
669            },
670            "dYdX orderbook unsubscription",
671        );
672
673        Ok(())
674    }
675
676    fn unsubscribe_quotes(&mut self, cmd: &UnsubscribeQuotes) -> anyhow::Result<()> {
677        log::debug!(
678            "unsubscribe_quotes for {}: removing quote subscription",
679            cmd.instrument_id
680        );
681
682        self.active_quote_subs.remove(&cmd.instrument_id);
683
684        let ws = self.ws_client.clone();
685        let instrument_id = cmd.instrument_id;
686
687        self.spawn_ws(
688            async move {
689                ws.unsubscribe_orderbook(instrument_id)
690                    .await
691                    .context("orderbook unsubscription (for quotes)")
692            },
693            "dYdX orderbook unsubscription (quotes)",
694        );
695
696        Ok(())
697    }
698
699    fn unsubscribe_trades(&mut self, cmd: &UnsubscribeTrades) -> anyhow::Result<()> {
700        self.active_trade_subs.remove(&cmd.instrument_id);
701
702        let ws = self.ws_client.clone();
703        let instrument_id = cmd.instrument_id;
704
705        self.spawn_ws(
706            async move {
707                ws.unsubscribe_trades(instrument_id)
708                    .await
709                    .context("trade unsubscription")
710            },
711            "dYdX trade unsubscription",
712        );
713
714        Ok(())
715    }
716
717    fn unsubscribe_mark_prices(&mut self, cmd: &UnsubscribeMarkPrices) -> anyhow::Result<()> {
718        self.active_mark_price_subs.remove(&cmd.instrument_id);
719        log::debug!("Unsubscribed from mark prices for {}", cmd.instrument_id);
720        Ok(())
721    }
722
723    fn unsubscribe_index_prices(&mut self, cmd: &UnsubscribeIndexPrices) -> anyhow::Result<()> {
724        self.active_index_price_subs.remove(&cmd.instrument_id);
725        log::debug!("Unsubscribed from index prices for {}", cmd.instrument_id);
726        Ok(())
727    }
728
729    fn unsubscribe_bars(&mut self, cmd: &UnsubscribeBars) -> anyhow::Result<()> {
730        let ws = self.ws_client.clone();
731        let instrument_id = cmd.bar_type.instrument_id();
732        let spec = cmd.bar_type.spec();
733
734        let resolution = Self::map_bar_spec_to_resolution(&spec)?;
735
736        self.active_bar_subs
737            .remove(&(instrument_id, resolution.to_string()));
738
739        let ticker = extract_raw_symbol(instrument_id.symbol.as_str());
740        let topic = format!("{ticker}/{resolution}");
741        self.bar_type_mappings.remove(&topic);
742
743        self.spawn_ws(
744            async move {
745                ws.unsubscribe_candles(instrument_id, resolution)
746                    .await
747                    .context("candles unsubscription")
748            },
749            "dYdX candles unsubscription",
750        );
751
752        Ok(())
753    }
754
755    fn unsubscribe_funding_rates(&mut self, cmd: &UnsubscribeFundingRates) -> anyhow::Result<()> {
756        self.active_funding_rate_subs.remove(&cmd.instrument_id);
757        log::debug!("Unsubscribed from funding rates for {}", cmd.instrument_id);
758        Ok(())
759    }
760
761    fn unsubscribe_instrument_status(
762        &mut self,
763        cmd: &UnsubscribeInstrumentStatus,
764    ) -> anyhow::Result<()> {
765        self.active_instrument_status_subs
766            .remove(&cmd.instrument_id);
767        log::debug!(
768            "Unsubscribed from instrument status for {}",
769            cmd.instrument_id
770        );
771        Ok(())
772    }
773
774    fn request_instrument(&self, request: RequestInstrument) -> anyhow::Result<()> {
775        if request.start.is_some() {
776            log::warn!(
777                "Requesting instrument {} with specified `start` which has no effect",
778                request.instrument_id
779            );
780        }
781
782        if request.end.is_some() {
783            log::warn!(
784                "Requesting instrument {} with specified `end` which has no effect",
785                request.instrument_id
786            );
787        }
788
789        let instrument_cache = self.instrument_cache.clone();
790        let sender = self.data_sender.clone();
791        let http = self.http_client.clone();
792        let instrument_id = request.instrument_id;
793        let request_id = request.request_id;
794        let client_id = request.client_id.unwrap_or(self.client_id);
795        let start = request.start;
796        let end = request.end;
797        let params = request.params;
798        let clock = self.clock;
799        let start_nanos = datetime_to_unix_nanos(start);
800        let end_nanos = datetime_to_unix_nanos(end);
801
802        self.spawn_command(async move {
803            let instrument = match http.request_instruments(None, None, None).await {
804                Ok(instruments) => {
805                    for inst in &instruments {
806                        instrument_cache.insert_instrument_only(inst.clone());
807                    }
808                    instruments.into_iter().find(|i| i.id() == instrument_id)
809                }
810                Err(e) => {
811                    log::error!("Failed to fetch instruments from dYdX: {e:?}");
812                    None
813                }
814            };
815
816            if let Some(inst) = instrument {
817                let response = DataResponse::Instrument(Box::new(InstrumentResponse::new(
818                    request_id,
819                    client_id,
820                    instrument_id,
821                    inst,
822                    start_nanos,
823                    end_nanos,
824                    clock.get_time_ns(),
825                    params,
826                )));
827
828                if let Err(e) = sender.send(DataEvent::Response(response)) {
829                    log::error!("Failed to send instrument response: {e}");
830                }
831            } else {
832                log::error!("Instrument {instrument_id} not found");
833            }
834        });
835
836        Ok(())
837    }
838
839    fn request_instruments(&self, request: RequestInstruments) -> anyhow::Result<()> {
840        let http = self.http_client.clone();
841        let sender = self.data_sender.clone();
842        let instrument_cache = self.instrument_cache.clone();
843        let request_id = request.request_id;
844        let client_id = request.client_id.unwrap_or(self.client_id);
845        let venue = self.venue();
846        let start = request.start;
847        let end = request.end;
848        let params = request.params;
849        let clock = self.clock;
850        let start_nanos = datetime_to_unix_nanos(start);
851        let end_nanos = datetime_to_unix_nanos(end);
852
853        self.spawn_command(async move {
854            match http.request_instruments(None, None, None).await {
855                Ok(instruments) => {
856                    log::debug!("Fetched {} instruments from dYdX", instruments.len());
857
858                    for instrument in &instruments {
859                        instrument_cache.insert_instrument_only(instrument.clone());
860                    }
861
862                    let response = DataResponse::Instruments(InstrumentsResponse::new(
863                        request_id,
864                        client_id,
865                        venue,
866                        instruments,
867                        start_nanos,
868                        end_nanos,
869                        clock.get_time_ns(),
870                        params,
871                    ));
872
873                    if let Err(e) = sender.send(DataEvent::Response(response)) {
874                        log::error!("Failed to send instruments response: {e}");
875                    }
876                }
877                Err(e) => {
878                    log::error!("Failed to fetch instruments from dYdX: {e:?}");
879
880                    let response = DataResponse::Instruments(InstrumentsResponse::new(
881                        request_id,
882                        client_id,
883                        venue,
884                        Vec::new(),
885                        start_nanos,
886                        end_nanos,
887                        clock.get_time_ns(),
888                        params,
889                    ));
890
891                    if let Err(e) = sender.send(DataEvent::Response(response)) {
892                        log::error!("Failed to send empty instruments response: {e}");
893                    }
894                }
895            }
896        });
897
898        Ok(())
899    }
900
901    fn request_book_snapshot(&self, request: RequestBookSnapshot) -> anyhow::Result<()> {
902        if request.depth.is_some() {
903            log::warn!(
904                "Requesting book snapshot for {} with specified `depth` which has no effect",
905                request.instrument_id
906            );
907        }
908
909        let http_client = self.http_client.clone();
910        let sender = self.data_sender.clone();
911        let instrument_id = request.instrument_id;
912        let request_id = request.request_id;
913        let client_id = request.client_id.unwrap_or(self.client_id);
914        let params = request.params;
915        let clock = self.clock;
916
917        self.spawn_command(async move {
918            let mut book = OrderBook::new(instrument_id, BookType::L2_MBP);
919
920            match http_client.request_orderbook_snapshot(instrument_id).await {
921                Ok(deltas) => {
922                    if let Err(e) = book.apply_deltas(&deltas) {
923                        log::error!("Failed to apply book snapshot for {instrument_id}: {e}");
924                        book.reset();
925                    }
926                }
927                Err(e) => {
928                    log::error!("Book snapshot request failed for {instrument_id}: {e:?}");
929                }
930            }
931
932            let response = DataResponse::Book(BookResponse::new(
933                request_id,
934                client_id,
935                instrument_id,
936                book,
937                None,
938                None,
939                clock.get_time_ns(),
940                params,
941            ));
942
943            if let Err(e) = sender.send(DataEvent::Response(response)) {
944                log::error!("Failed to send book snapshot response: {e}");
945            }
946        });
947
948        Ok(())
949    }
950
951    fn request_trades(&self, request: RequestTrades) -> anyhow::Result<()> {
952        let http_client = self.http_client.clone();
953        let sender = self.data_sender.clone();
954        let instrument_id = request.instrument_id;
955        let start = request.start;
956        let end = request.end;
957        let limit = request.limit.map(|n| n.get() as u32);
958        let request_id = request.request_id;
959        let client_id = request.client_id.unwrap_or(self.client_id);
960        let params = request.params;
961        let clock = self.clock;
962        let start_nanos = datetime_to_unix_nanos(start);
963        let end_nanos = datetime_to_unix_nanos(end);
964
965        self.spawn_command(async move {
966            match http_client
967                .request_trade_ticks(instrument_id, start, end, limit)
968                .await
969                .context("failed to request trades from dYdX")
970            {
971                Ok(trades) => {
972                    let response = DataResponse::Trades(TradesResponse::new(
973                        request_id,
974                        client_id,
975                        instrument_id,
976                        trades,
977                        start_nanos,
978                        end_nanos,
979                        clock.get_time_ns(),
980                        params,
981                    ));
982
983                    if let Err(e) = sender.send(DataEvent::Response(response)) {
984                        log::error!("Failed to send trades response: {e}");
985                    }
986                }
987                Err(e) => {
988                    log::error!("Trade request failed for {instrument_id}: {e:?}");
989
990                    let response = DataResponse::Trades(TradesResponse::new(
991                        request_id,
992                        client_id,
993                        instrument_id,
994                        Vec::new(),
995                        start_nanos,
996                        end_nanos,
997                        clock.get_time_ns(),
998                        params,
999                    ));
1000
1001                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1002                        log::error!("Failed to send empty trades response: {e}");
1003                    }
1004                }
1005            }
1006        });
1007
1008        Ok(())
1009    }
1010
1011    fn request_bars(&self, request: RequestBars) -> anyhow::Result<()> {
1012        let http_client = self.http_client.clone();
1013        let sender = self.data_sender.clone();
1014        let bar_type = request.bar_type;
1015        let start = request.start;
1016        let end = request.end;
1017        let limit = request.limit.map(|n| n.get() as u32);
1018        let request_id = request.request_id;
1019        let client_id = request.client_id.unwrap_or(self.client_id);
1020        let params = request.params;
1021        let clock = self.clock;
1022        let start_nanos = datetime_to_unix_nanos(start);
1023        let end_nanos = datetime_to_unix_nanos(end);
1024
1025        self.spawn_command(async move {
1026            match http_client
1027                .request_bars(bar_type, start, end, limit, true)
1028                .await
1029                .context("failed to request bars from dYdX")
1030            {
1031                Ok(bars) => {
1032                    let response = DataResponse::Bars(BarsResponse::new(
1033                        request_id,
1034                        client_id,
1035                        bar_type,
1036                        bars,
1037                        start_nanos,
1038                        end_nanos,
1039                        clock.get_time_ns(),
1040                        params,
1041                    ));
1042
1043                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1044                        log::error!("Failed to send bars response: {e}");
1045                    }
1046                }
1047                Err(e) => {
1048                    log::error!("Bar request failed for {bar_type}: {e:?}");
1049
1050                    let response = DataResponse::Bars(BarsResponse::new(
1051                        request_id,
1052                        client_id,
1053                        bar_type,
1054                        Vec::new(),
1055                        start_nanos,
1056                        end_nanos,
1057                        clock.get_time_ns(),
1058                        params,
1059                    ));
1060
1061                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1062                        log::error!("Failed to send empty bars response: {e}");
1063                    }
1064                }
1065            }
1066        });
1067
1068        Ok(())
1069    }
1070
1071    fn request_funding_rates(&self, request: RequestFundingRates) -> anyhow::Result<()> {
1072        let http_client = self.http_client.clone();
1073        let sender = self.data_sender.clone();
1074        let instrument_id = request.instrument_id;
1075        let start = request.start;
1076        let end = request.end;
1077        let limit = request.limit.map(|n| n.get() as u32);
1078        let request_id = request.request_id;
1079        let client_id = request.client_id.unwrap_or(self.client_id);
1080        let params = request.params;
1081        let clock = self.clock;
1082        let start_nanos = datetime_to_unix_nanos(start);
1083        let end_nanos = datetime_to_unix_nanos(end);
1084
1085        self.spawn_command(async move {
1086            match http_client
1087                .request_funding_rates(instrument_id, start, end, limit)
1088                .await
1089                .context("failed to request funding rates from dYdX")
1090            {
1091                Ok(funding_rates) => {
1092                    let response = DataResponse::FundingRates(FundingRatesResponse::new(
1093                        request_id,
1094                        client_id,
1095                        instrument_id,
1096                        funding_rates,
1097                        start_nanos,
1098                        end_nanos,
1099                        clock.get_time_ns(),
1100                        params,
1101                    ));
1102
1103                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1104                        log::error!("Failed to send funding rates response: {e}");
1105                    }
1106                }
1107                Err(e) => {
1108                    log::error!("Funding rates request failed for {instrument_id}: {e:?}");
1109
1110                    let response = DataResponse::FundingRates(FundingRatesResponse::new(
1111                        request_id,
1112                        client_id,
1113                        instrument_id,
1114                        Vec::new(),
1115                        start_nanos,
1116                        end_nanos,
1117                        clock.get_time_ns(),
1118                        params,
1119                    ));
1120
1121                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1122                        log::error!("Failed to send empty funding rates response: {e}");
1123                    }
1124                }
1125            }
1126        });
1127
1128        Ok(())
1129    }
1130}
1131
1132impl DydxDataClient {
1133    /// Returns a cached instrument by InstrumentId.
1134    #[must_use]
1135    pub fn get_instrument(&self, instrument_id: &InstrumentId) -> Option<InstrumentAny> {
1136        self.instrument_cache.get(instrument_id)
1137    }
1138
1139    /// Returns all cached instruments.
1140    #[must_use]
1141    pub fn get_instruments(&self) -> Vec<InstrumentAny> {
1142        self.instrument_cache.all_instruments()
1143    }
1144
1145    /// Caches a single instrument.
1146    pub fn cache_instrument(&self, instrument: InstrumentAny) {
1147        self.instrument_cache.insert_instrument_only(instrument);
1148    }
1149
1150    /// Caches multiple instruments.
1151    ///
1152    /// Clears the existing cache first, then adds all provided instruments.
1153    pub fn cache_instruments(&self, instruments: Vec<InstrumentAny>) {
1154        self.instrument_cache.clear();
1155        self.instrument_cache.insert_instruments_only(instruments);
1156    }
1157
1158    fn ensure_order_book(&self, instrument_id: InstrumentId, book_type: BookType) {
1159        self.order_books
1160            .entry(instrument_id)
1161            .or_insert_with(|| OrderBook::new(instrument_id, book_type));
1162    }
1163
1164    /// Returns the BarType for a given WebSocket candle topic.
1165    #[must_use]
1166    pub fn get_bar_type_for_topic(&self, topic: &str) -> Option<BarType> {
1167        self.bar_type_mappings.load().get(topic).copied()
1168    }
1169
1170    /// Returns all registered bar topics.
1171    #[must_use]
1172    pub fn get_bar_topics(&self) -> Vec<String> {
1173        self.bar_type_mappings.load().keys().cloned().collect()
1174    }
1175
1176    fn handle_ws_message(message: DydxWsOutputMessage, ctx: &WsMessageContext) {
1177        let ts_init = ctx.clock.get_time_ns();
1178
1179        match message {
1180            DydxWsOutputMessage::Trades { id, contents } => {
1181                let Some(instrument) = ctx.instrument_cache.get_by_market(&id) else {
1182                    log::warn!("No instrument cached for market {id}");
1183                    return;
1184                };
1185                let instrument_id = instrument.id();
1186
1187                match ws_parse::parse_trade_ticks(instrument_id, &instrument, &contents, ts_init) {
1188                    Ok(data) => {
1189                        Self::handle_data_message(
1190                            data,
1191                            &ctx.data_sender,
1192                            &ctx.incomplete_bars,
1193                            ctx.clock,
1194                        );
1195                    }
1196                    Err(e) => log::error!("Failed to parse trade ticks for {id}: {e}"),
1197                }
1198            }
1199            DydxWsOutputMessage::OrderbookSnapshot { id, contents } => {
1200                let Some(instrument) = ctx.instrument_cache.get_by_market(&id) else {
1201                    log::warn!("No instrument cached for market {id}");
1202                    return;
1203                };
1204                let instrument_id = instrument.id();
1205
1206                match ws_parse::parse_orderbook_snapshot(
1207                    &instrument_id,
1208                    &contents,
1209                    instrument.price_precision(),
1210                    instrument.size_precision(),
1211                    ts_init,
1212                ) {
1213                    Ok(deltas) => {
1214                        Self::handle_deltas_message(
1215                            deltas,
1216                            &ctx.data_sender,
1217                            &ctx.order_books,
1218                            &ctx.last_quotes,
1219                            &ctx.instrument_cache,
1220                            &ctx.active_quote_subs,
1221                            &ctx.active_delta_subs,
1222                        );
1223                    }
1224                    Err(e) => log::error!("Failed to parse orderbook snapshot for {id}: {e}"),
1225                }
1226            }
1227            DydxWsOutputMessage::OrderbookUpdate { id, contents } => {
1228                let Some(instrument) = ctx.instrument_cache.get_by_market(&id) else {
1229                    log::warn!("No instrument cached for market {id}");
1230                    return;
1231                };
1232                let instrument_id = instrument.id();
1233
1234                match ws_parse::parse_orderbook_deltas(
1235                    &instrument_id,
1236                    &contents,
1237                    instrument.price_precision(),
1238                    instrument.size_precision(),
1239                    ts_init,
1240                ) {
1241                    Ok(deltas) => {
1242                        Self::handle_deltas_message(
1243                            deltas,
1244                            &ctx.data_sender,
1245                            &ctx.order_books,
1246                            &ctx.last_quotes,
1247                            &ctx.instrument_cache,
1248                            &ctx.active_quote_subs,
1249                            &ctx.active_delta_subs,
1250                        );
1251                    }
1252                    Err(e) => log::error!("Failed to parse orderbook deltas for {id}: {e}"),
1253                }
1254            }
1255            DydxWsOutputMessage::OrderbookBatch { id, updates } => {
1256                let Some(instrument) = ctx.instrument_cache.get_by_market(&id) else {
1257                    log::warn!("No instrument cached for market {id}");
1258                    return;
1259                };
1260                let instrument_id = instrument.id();
1261                let price_precision = instrument.price_precision();
1262                let size_precision = instrument.size_precision();
1263
1264                let mut all_deltas = Vec::new();
1265                let last_idx = updates.len().saturating_sub(1);
1266
1267                for (i, update) in updates.iter().enumerate() {
1268                    let is_last = i == last_idx;
1269                    let result = if is_last {
1270                        ws_parse::parse_orderbook_deltas(
1271                            &instrument_id,
1272                            update,
1273                            price_precision,
1274                            size_precision,
1275                            ts_init,
1276                        )
1277                        .map(|d| d.deltas)
1278                    } else {
1279                        ws_parse::parse_orderbook_deltas_with_flag(
1280                            &instrument_id,
1281                            update,
1282                            price_precision,
1283                            size_precision,
1284                            ts_init,
1285                            false,
1286                        )
1287                    };
1288
1289                    match result {
1290                        Ok(deltas) => all_deltas.extend(deltas),
1291                        Err(e) => {
1292                            log::error!("Failed to parse orderbook batch delta {i} for {id}: {e}");
1293                            return;
1294                        }
1295                    }
1296                }
1297
1298                if all_deltas.is_empty() {
1299                    return;
1300                }
1301                let deltas = OrderBookDeltas::new(instrument_id, all_deltas);
1302                Self::handle_deltas_message(
1303                    deltas,
1304                    &ctx.data_sender,
1305                    &ctx.order_books,
1306                    &ctx.last_quotes,
1307                    &ctx.instrument_cache,
1308                    &ctx.active_quote_subs,
1309                    &ctx.active_delta_subs,
1310                );
1311            }
1312            DydxWsOutputMessage::Candles { id, contents } => {
1313                let parts: Vec<&str> = id.splitn(2, '/').collect();
1314                if parts.len() != 2 {
1315                    log::warn!("Unexpected candle topic format: {id}");
1316                    return;
1317                }
1318                let ticker = parts[0];
1319
1320                let Some(bar_type) = ctx.bar_type_mappings.load().get(&id).copied() else {
1321                    log::debug!("No bar type mapping for candle topic {id}");
1322                    return;
1323                };
1324
1325                let Some(instrument) = ctx.instrument_cache.get_by_market(ticker) else {
1326                    log::warn!("No instrument cached for market {ticker}");
1327                    return;
1328                };
1329
1330                match ws_parse::parse_candle_bar(
1331                    bar_type,
1332                    &instrument,
1333                    &contents,
1334                    ctx.bars_timestamp_on_close,
1335                    ts_init,
1336                ) {
1337                    Ok(bar) => {
1338                        let prev = ctx.pending_bars.get(&id).map(|r| *r);
1339                        if let Some(prev_bar) = prev
1340                            && bar.ts_event != prev_bar.ts_event
1341                        {
1342                            Self::emit_bar_guarded(prev_bar, ctx);
1343                        }
1344                        ctx.pending_bars.insert(id, bar);
1345                    }
1346                    Err(e) => log::error!("Failed to parse candle bar for {id}: {e}"),
1347                }
1348            }
1349            DydxWsOutputMessage::Markets(contents) => {
1350                Self::handle_markets_message(&contents, ctx, ts_init);
1351            }
1352            DydxWsOutputMessage::SubaccountSubscribed(_) => {
1353                log::debug!("Ignoring subaccount subscribed on data client");
1354            }
1355            DydxWsOutputMessage::SubaccountsChannelData(_) => {
1356                log::debug!("Ignoring subaccounts channel data on data client");
1357            }
1358            DydxWsOutputMessage::BlockHeight { .. } => {
1359                log::debug!("Ignoring block height on data client");
1360            }
1361            DydxWsOutputMessage::Error(err) => {
1362                log::warn!("dYdX WS error: {err}");
1363            }
1364            DydxWsOutputMessage::Reconnected { topics } => {
1365                let reconnected_candles = candle_ids_from_topics(&topics);
1366                ctx.pending_bars
1367                    .retain(|id, _| !reconnected_candles.contains(id));
1368
1369                let total_subs = ctx.active_quote_subs.len()
1370                    + ctx.active_delta_subs.len()
1371                    + ctx.active_trade_subs.len()
1372                    + ctx.active_bar_subs.len();
1373
1374                log::info!(
1375                    "dYdX WS reconnected; handler replayed channel subscriptions (active data subscriptions: total={}, quotes={}, deltas={}, trades={}, bars={})",
1376                    total_subs,
1377                    ctx.active_quote_subs.len(),
1378                    ctx.active_delta_subs.len(),
1379                    ctx.active_trade_subs.len(),
1380                    ctx.active_bar_subs.len()
1381                );
1382            }
1383        }
1384    }
1385
1386    fn instrument_id_from_ticker(ticker: &str) -> InstrumentId {
1387        let symbol = format!("{ticker}-PERP");
1388        InstrumentId::new(Symbol::new(&symbol), *DYDX_VENUE)
1389    }
1390
1391    fn handle_markets_message(
1392        contents: &crate::websocket::messages::DydxMarketsContents,
1393        ctx: &WsMessageContext,
1394        ts_init: nautilus_core::UnixNanos,
1395    ) {
1396        if let Some(ref oracle_prices) = contents.oracle_prices {
1397            for (ticker, oracle_data) in oracle_prices {
1398                let instrument_id = Self::instrument_id_from_ticker(ticker);
1399
1400                let Ok(price) = parse_price(&oracle_data.oracle_price, "oracle_price") else {
1401                    log::warn!("Failed to parse oracle price for {ticker}");
1402                    continue;
1403                };
1404
1405                if ctx.active_mark_price_subs.contains(&instrument_id) {
1406                    let mark_price = MarkPriceUpdate::new(instrument_id, price, ts_init, ts_init);
1407                    let data = NautilusData::MarkPrice(mark_price);
1408                    if let Err(e) = ctx.data_sender.send(DataEvent::Data(data)) {
1409                        log::error!("Failed to emit mark price for {instrument_id}: {e}");
1410                    }
1411                }
1412
1413                if ctx.active_index_price_subs.contains(&instrument_id) {
1414                    let index_price = IndexPriceUpdate::new(instrument_id, price, ts_init, ts_init);
1415                    let data = NautilusData::IndexPrice(index_price);
1416                    if let Err(e) = ctx.data_sender.send(DataEvent::Data(data)) {
1417                        log::error!("Failed to emit index price for {instrument_id}: {e}");
1418                    }
1419                }
1420            }
1421        }
1422
1423        Self::handle_markets_trading_data(contents.trading.as_ref(), ctx, ts_init, false);
1424        Self::handle_markets_trading_data(contents.markets.as_ref(), ctx, ts_init, true);
1425    }
1426
1427    fn handle_markets_trading_data(
1428        trading: Option<
1429            &std::collections::HashMap<String, crate::websocket::messages::DydxMarketTradingUpdate>,
1430        >,
1431        ctx: &WsMessageContext,
1432        ts_init: nautilus_core::UnixNanos,
1433        is_snapshot: bool,
1434    ) {
1435        let Some(trading_map) = trading else {
1436            return;
1437        };
1438
1439        for (ticker, update) in trading_map {
1440            let instrument_id = Self::instrument_id_from_ticker(ticker);
1441
1442            if let Some(status) = &update.status {
1443                if *status == crate::common::enums::DydxMarketStatus::Unknown {
1444                    log::warn!("Skipping unmodeled dYdX market status for {instrument_id}");
1445                } else {
1446                    let action = MarketStatusAction::from(*status);
1447                    let is_trading =
1448                        matches!(status, crate::common::enums::DydxMarketStatus::Active);
1449
1450                    let instrument_status = InstrumentStatus::new(
1451                        instrument_id,
1452                        action,
1453                        ts_init,
1454                        ts_init,
1455                        None,
1456                        None,
1457                        Some(is_trading),
1458                        None,
1459                        None,
1460                    );
1461
1462                    ctx.last_instrument_statuses
1463                        .insert(instrument_id, instrument_status);
1464
1465                    if ctx.active_instrument_status_subs.contains(&instrument_id)
1466                        && let Err(e) = ctx
1467                            .data_sender
1468                            .send(DataEvent::InstrumentStatus(instrument_status))
1469                    {
1470                        log::error!("Failed to emit instrument status for {instrument_id}: {e}");
1471                    }
1472                }
1473            }
1474
1475            let ticker_ustr = Ustr::from(ticker.as_str());
1476            if !ctx.seen_tickers.contains(&ticker_ustr) {
1477                let is_active = update
1478                    .status
1479                    .as_ref()
1480                    .is_none_or(|s| matches!(s, crate::common::enums::DydxMarketStatus::Active));
1481
1482                if ctx.instrument_cache.get_by_market(ticker).is_some() {
1483                    ctx.seen_tickers.insert(ticker_ustr);
1484                } else if is_active {
1485                    ctx.seen_tickers.insert(ticker_ustr);
1486                    Self::handle_new_instrument_discovered(ticker, ctx);
1487                }
1488            }
1489
1490            if let Some(ref rate_str) = update.next_funding_rate {
1491                if let Ok(rate) = Decimal::from_str(rate_str) {
1492                    if ctx.active_funding_rate_subs.contains(&instrument_id) {
1493                        let funding_rate = FundingRateUpdate {
1494                            instrument_id,
1495                            rate,
1496                            interval: Some(60),
1497                            next_funding_ns: None,
1498                            ts_event: ts_init,
1499                            ts_init,
1500                        };
1501
1502                        if let Err(e) = ctx.data_sender.send(DataEvent::FundingRate(funding_rate)) {
1503                            log::error!("Failed to emit funding rate for {instrument_id}: {e}");
1504                        }
1505                    }
1506                } else {
1507                    log::warn!("Failed to parse next_funding_rate for {ticker}: {rate_str}");
1508                }
1509            }
1510
1511            if is_snapshot
1512                && let Some(ref oracle_price_str) = update.oracle_price
1513                && let Ok(price) = parse_price(oracle_price_str, "oracle_price")
1514            {
1515                if ctx.active_mark_price_subs.contains(&instrument_id) {
1516                    let mark_price = MarkPriceUpdate::new(instrument_id, price, ts_init, ts_init);
1517                    let data = NautilusData::MarkPrice(mark_price);
1518
1519                    if let Err(e) = ctx.data_sender.send(DataEvent::Data(data)) {
1520                        log::error!("Failed to emit mark price for {instrument_id}: {e}");
1521                    }
1522                }
1523
1524                if ctx.active_index_price_subs.contains(&instrument_id) {
1525                    let index_price = IndexPriceUpdate::new(instrument_id, price, ts_init, ts_init);
1526                    let data = NautilusData::IndexPrice(index_price);
1527
1528                    if let Err(e) = ctx.data_sender.send(DataEvent::Data(data)) {
1529                        log::error!("Failed to emit index price for {instrument_id}: {e}");
1530                    }
1531                }
1532            }
1533        }
1534    }
1535
1536    fn emit_bar_guarded(bar: Bar, ctx: &WsMessageContext) {
1537        let current_time_ns = ctx.clock.get_time_ns();
1538        if bar.ts_event <= current_time_ns {
1539            ctx.incomplete_bars.remove(&bar.bar_type);
1540            if let Err(e) = ctx
1541                .data_sender
1542                .send(DataEvent::Data(NautilusData::Bar(bar)))
1543            {
1544                log::error!("Failed to emit completed bar: {e}");
1545            }
1546        } else {
1547            ctx.incomplete_bars.insert(bar.bar_type, bar);
1548        }
1549    }
1550
1551    fn handle_new_instrument_discovered(ticker: &str, ctx: &WsMessageContext) {
1552        log::debug!("New instrument discovered via WebSocket: {ticker}");
1553
1554        let http_client = ctx.http_client.clone();
1555        let ws_client = ctx.ws_client.clone();
1556        let data_sender = ctx.data_sender.clone();
1557        let ticker = ticker.to_string();
1558
1559        if let Err(e) = ctx.command_spawner.spawn(async move {
1560            match http_client.fetch_and_cache_single_instrument(&ticker).await {
1561                Ok(Some(instrument)) => {
1562                    ws_client.cache_instrument(instrument.clone());
1563                    if let Err(e) = data_sender.send(DataEvent::Instrument(instrument)) {
1564                        log::error!("Failed to emit new instrument: {e}");
1565                    }
1566                    log::debug!("Fetched and cached new instrument: {ticker}");
1567                }
1568                Ok(None) => {
1569                    log::warn!("New instrument {ticker} not found or inactive");
1570                }
1571                Err(e) => {
1572                    log::error!("Failed to fetch new instrument {ticker}: {e}");
1573                }
1574            }
1575        }) {
1576            log::warn!("Skipping new dYdX instrument fetch after shutdown began: {e}");
1577        }
1578    }
1579
1580    fn handle_data_message(
1581        payloads: Vec<NautilusData>,
1582        data_sender: &EventSender<DataEvent>,
1583        incomplete_bars: &Arc<DashMap<BarType, Bar>>,
1584        clock: &'static AtomicTime,
1585    ) {
1586        for data in payloads {
1587            if let NautilusData::Bar(bar) = data {
1588                Self::handle_bar_message(bar, data_sender, incomplete_bars, clock);
1589            } else if let Err(e) = data_sender.send(DataEvent::Data(data)) {
1590                log::error!("Failed to emit data event: {e}");
1591            }
1592        }
1593    }
1594
1595    fn handle_bar_message(
1596        bar: Bar,
1597        data_sender: &EventSender<DataEvent>,
1598        incomplete_bars: &Arc<DashMap<BarType, Bar>>,
1599        clock: &'static AtomicTime,
1600    ) {
1601        let current_time_ns = clock.get_time_ns();
1602        let bar_type = bar.bar_type;
1603
1604        if bar.ts_event <= current_time_ns {
1605            incomplete_bars.remove(&bar_type);
1606
1607            if let Err(e) = data_sender.send(DataEvent::Data(NautilusData::Bar(bar))) {
1608                log::error!("Failed to emit completed bar: {e}");
1609            }
1610        } else {
1611            log::trace!(
1612                "Caching incomplete bar for {} (ts_event={}, current={})",
1613                bar_type,
1614                bar.ts_event,
1615                current_time_ns
1616            );
1617            incomplete_bars.insert(bar_type, bar);
1618        }
1619    }
1620
1621    fn resolve_crossed_order_book(
1622        book: &mut OrderBook,
1623        venue_deltas: &OrderBookDeltas,
1624        instrument: &InstrumentAny,
1625    ) -> anyhow::Result<OrderBookDeltas> {
1626        let instrument_id = venue_deltas.instrument_id;
1627        let ts_init = venue_deltas.ts_init;
1628        let mut all_deltas = venue_deltas.deltas.clone();
1629
1630        // If the input batch is a snapshot, every synthetic and terminator delta must
1631        // carry F_SNAPSHOT as well so consumers apply the whole batch as one
1632        // replacement image rather than a snapshot followed by standalone updates.
1633        let snapshot_flag = RecordFlag::F_SNAPSHOT as u8;
1634        let is_snapshot_batch = venue_deltas
1635            .deltas
1636            .iter()
1637            .any(|d| d.flags & snapshot_flag != 0);
1638        let synthetic_flags = if is_snapshot_batch { snapshot_flag } else { 0 };
1639
1640        book.apply_deltas(venue_deltas)?;
1641
1642        let mut is_crossed = if let (Some(bid_price), Some(ask_price)) =
1643            (book.best_bid_price(), book.best_ask_price())
1644        {
1645            bid_price >= ask_price
1646        } else {
1647            false
1648        };
1649
1650        while is_crossed {
1651            log::debug!(
1652                "Resolving crossed order book for {}: bid={:?} >= ask={:?}",
1653                instrument_id,
1654                book.best_bid_price(),
1655                book.best_ask_price()
1656            );
1657
1658            let bid_price = match book.best_bid_price() {
1659                Some(p) => p,
1660                None => break,
1661            };
1662            let ask_price = match book.best_ask_price() {
1663                Some(p) => p,
1664                None => break,
1665            };
1666            let bid_size = match book.best_bid_size() {
1667                Some(s) => s,
1668                None => break,
1669            };
1670            let ask_size = match book.best_ask_size() {
1671                Some(s) => s,
1672                None => break,
1673            };
1674
1675            let mut temp_deltas = Vec::new();
1676
1677            if bid_size > ask_size {
1678                let new_bid_size = Quantity::from_decimal_dp(
1679                    bid_size.as_decimal() - ask_size.as_decimal(),
1680                    instrument.size_precision(),
1681                )?;
1682                temp_deltas.push(OrderBookDelta::new(
1683                    instrument_id,
1684                    BookAction::Update,
1685                    BookOrder::new(OrderSide::Buy, bid_price, new_bid_size, 0),
1686                    synthetic_flags,
1687                    0,
1688                    ts_init,
1689                    ts_init,
1690                ));
1691                temp_deltas.push(OrderBookDelta::new(
1692                    instrument_id,
1693                    BookAction::Delete,
1694                    BookOrder::new(
1695                        OrderSide::Sell,
1696                        ask_price,
1697                        Quantity::zero(instrument.size_precision()),
1698                        0,
1699                    ),
1700                    synthetic_flags,
1701                    0,
1702                    ts_init,
1703                    ts_init,
1704                ));
1705            } else if bid_size < ask_size {
1706                let new_ask_size = Quantity::from_decimal_dp(
1707                    ask_size.as_decimal() - bid_size.as_decimal(),
1708                    instrument.size_precision(),
1709                )?;
1710                temp_deltas.push(OrderBookDelta::new(
1711                    instrument_id,
1712                    BookAction::Update,
1713                    BookOrder::new(OrderSide::Sell, ask_price, new_ask_size, 0),
1714                    synthetic_flags,
1715                    0,
1716                    ts_init,
1717                    ts_init,
1718                ));
1719                temp_deltas.push(OrderBookDelta::new(
1720                    instrument_id,
1721                    BookAction::Delete,
1722                    BookOrder::new(
1723                        OrderSide::Buy,
1724                        bid_price,
1725                        Quantity::zero(instrument.size_precision()),
1726                        0,
1727                    ),
1728                    synthetic_flags,
1729                    0,
1730                    ts_init,
1731                    ts_init,
1732                ));
1733            } else {
1734                temp_deltas.push(OrderBookDelta::new(
1735                    instrument_id,
1736                    BookAction::Delete,
1737                    BookOrder::new(
1738                        OrderSide::Buy,
1739                        bid_price,
1740                        Quantity::zero(instrument.size_precision()),
1741                        0,
1742                    ),
1743                    synthetic_flags,
1744                    0,
1745                    ts_init,
1746                    ts_init,
1747                ));
1748                temp_deltas.push(OrderBookDelta::new(
1749                    instrument_id,
1750                    BookAction::Delete,
1751                    BookOrder::new(
1752                        OrderSide::Sell,
1753                        ask_price,
1754                        Quantity::zero(instrument.size_precision()),
1755                        0,
1756                    ),
1757                    synthetic_flags,
1758                    0,
1759                    ts_init,
1760                    ts_init,
1761                ));
1762            }
1763
1764            let temp_deltas_obj = OrderBookDeltas::new(instrument_id, temp_deltas.clone());
1765            book.apply_deltas(&temp_deltas_obj)?;
1766            all_deltas.extend(temp_deltas);
1767
1768            is_crossed = if let (Some(bid_price), Some(ask_price)) =
1769                (book.best_bid_price(), book.best_ask_price())
1770            {
1771                bid_price >= ask_price
1772            } else {
1773                false
1774            };
1775        }
1776
1777        // Set F_LAST on the final delta, preserving F_SNAPSHOT when the batch is a
1778        // snapshot so consumers close the replacement image correctly.
1779        if let Some(last_delta) = all_deltas.last_mut() {
1780            last_delta.flags = synthetic_flags | RecordFlag::F_LAST as u8;
1781        }
1782
1783        Ok(OrderBookDeltas::new(instrument_id, all_deltas))
1784    }
1785
1786    fn handle_deltas_message(
1787        deltas: OrderBookDeltas,
1788        data_sender: &EventSender<DataEvent>,
1789        order_books: &Arc<DashMap<InstrumentId, OrderBook>>,
1790        last_quotes: &Arc<DashMap<InstrumentId, QuoteTick>>,
1791        instrument_cache: &Arc<InstrumentCache>,
1792        active_quote_subs: &Arc<AtomicSet<InstrumentId>>,
1793        active_delta_subs: &Arc<AtomicSet<InstrumentId>>,
1794    ) {
1795        let instrument_id = deltas.instrument_id;
1796
1797        let instrument = match instrument_cache.get(&instrument_id) {
1798            Some(inst) => inst,
1799            None => {
1800                log::error!("Cannot resolve crossed order book: no instrument for {instrument_id}");
1801                if active_delta_subs.contains(&instrument_id)
1802                    && let Err(e) = data_sender.send(DataEvent::Data(NautilusData::from(deltas)))
1803                {
1804                    log::error!("Failed to emit order book deltas: {e}");
1805                }
1806                return;
1807            }
1808        };
1809
1810        // Always maintain local orderbook -- both subscription types need book state
1811        let mut book = order_books
1812            .entry(instrument_id)
1813            .or_insert_with(|| OrderBook::new(instrument_id, BookType::L2_MBP));
1814
1815        let resolved_deltas =
1816            match Self::resolve_crossed_order_book(&mut book, &deltas, &instrument) {
1817                Ok(d) => d,
1818                Err(e) => {
1819                    log::error!("Failed to resolve crossed order book for {instrument_id}: {e}");
1820                    return;
1821                }
1822            };
1823
1824        if active_quote_subs.contains(&instrument_id) {
1825            // Edge case: If orderbook is empty after deltas, fall back to last quote
1826            let quote_opt = if let (Some(bid_price), Some(ask_price)) =
1827                (book.best_bid_price(), book.best_ask_price())
1828                && let (Some(bid_size), Some(ask_size)) =
1829                    (book.best_bid_size(), book.best_ask_size())
1830            {
1831                Some(QuoteTick::new(
1832                    instrument_id,
1833                    bid_price,
1834                    ask_price,
1835                    bid_size,
1836                    ask_size,
1837                    resolved_deltas.ts_event,
1838                    resolved_deltas.ts_init,
1839                ))
1840            } else if book.best_bid_price().is_none() && book.best_ask_price().is_none() {
1841                log::debug!(
1842                    "Empty orderbook for {instrument_id} after applying deltas, using last quote"
1843                );
1844                last_quotes.get(&instrument_id).map(|q| *q)
1845            } else {
1846                None
1847            };
1848
1849            if let Some(quote) = quote_opt {
1850                let emit_quote = !matches!(
1851                    last_quotes.get(&instrument_id),
1852                    Some(existing) if *existing == quote
1853                );
1854
1855                if emit_quote {
1856                    last_quotes.insert(instrument_id, quote);
1857                    if let Err(e) = data_sender.send(DataEvent::Data(NautilusData::Quote(quote))) {
1858                        log::error!("Failed to emit quote tick: {e}");
1859                    }
1860                }
1861            } else if book.best_bid_price().is_some() || book.best_ask_price().is_some() {
1862                log::debug!(
1863                    "Incomplete top-of-book for {instrument_id} (bid={:?}, ask={:?})",
1864                    book.best_bid_price(),
1865                    book.best_ask_price()
1866                );
1867            }
1868        }
1869
1870        if active_delta_subs.contains(&instrument_id) {
1871            let data: NautilusData = resolved_deltas.into();
1872            if let Err(e) = data_sender.send(DataEvent::Data(data)) {
1873                log::error!("Failed to emit order book deltas event: {e}");
1874            }
1875        }
1876    }
1877}
1878
1879struct WsMessageContext {
1880    clock: &'static AtomicTime,
1881    data_sender: EventSender<DataEvent>,
1882    instrument_cache: Arc<InstrumentCache>,
1883    order_books: Arc<DashMap<InstrumentId, OrderBook>>,
1884    last_quotes: Arc<DashMap<InstrumentId, QuoteTick>>,
1885    ws_client: DydxWebSocketClient,
1886    http_client: DydxHttpClient,
1887    active_quote_subs: Arc<AtomicSet<InstrumentId>>,
1888    active_delta_subs: Arc<AtomicSet<InstrumentId>>,
1889    active_trade_subs: Arc<AtomicSet<InstrumentId>>,
1890    active_bar_subs: Arc<AtomicMap<(InstrumentId, String), BarType>>,
1891    incomplete_bars: Arc<DashMap<BarType, Bar>>,
1892    bar_type_mappings: Arc<AtomicMap<String, BarType>>,
1893    active_mark_price_subs: Arc<AtomicSet<InstrumentId>>,
1894    active_index_price_subs: Arc<AtomicSet<InstrumentId>>,
1895    active_funding_rate_subs: Arc<AtomicSet<InstrumentId>>,
1896    active_instrument_status_subs: Arc<AtomicSet<InstrumentId>>,
1897    last_instrument_statuses: Arc<DashMap<InstrumentId, InstrumentStatus>>,
1898    bars_timestamp_on_close: bool,
1899    pending_bars: Arc<DashMap<String, Bar>>,
1900    seen_tickers: Arc<AtomicSet<Ustr>>,
1901    command_spawner: TaskSpawner,
1902}
1903
1904#[cfg(test)]
1905mod tests {
1906    use nautilus_core::UnixNanos;
1907    use nautilus_model::{
1908        data::{BookOrder, OrderBookDelta, OrderBookDeltas},
1909        enums::{BookAction, BookType, OrderSide, RecordFlag},
1910        identifiers::{InstrumentId, Symbol},
1911        instruments::{CryptoPerpetual, InstrumentAny},
1912        orderbook::OrderBook,
1913        types::{Currency, Price, Quantity},
1914    };
1915    use rstest::rstest;
1916    use rust_decimal_macros::dec;
1917
1918    use super::*;
1919    use crate::common::consts::DYDX_VENUE;
1920
1921    fn test_instrument() -> InstrumentAny {
1922        let instrument_id = InstrumentId::new(Symbol::new("BTC-USD-PERP"), *DYDX_VENUE);
1923        InstrumentAny::CryptoPerpetual(
1924            CryptoPerpetual::builder()
1925                .instrument_id(instrument_id)
1926                .raw_symbol(instrument_id.symbol)
1927                .base_currency(Currency::BTC())
1928                .quote_currency(Currency::USD())
1929                .settlement_currency(Currency::USD())
1930                .is_inverse(false)
1931                .price_precision(2)
1932                .size_precision(8)
1933                .price_increment(Price::new(0.01, 2))
1934                .size_increment(Quantity::new(0.00000001, 8))
1935                .ts_event(UnixNanos::default())
1936                .ts_init(UnixNanos::default())
1937                .build()
1938                .unwrap(),
1939        )
1940    }
1941
1942    fn seed_book_with_levels(
1943        instrument_id: InstrumentId,
1944        bids: &[(f64, f64)],
1945        asks: &[(f64, f64)],
1946    ) -> OrderBook {
1947        let mut book = OrderBook::new(instrument_id, BookType::L2_MBP);
1948        let ts = UnixNanos::default();
1949
1950        let mut deltas: Vec<OrderBookDelta> = Vec::new();
1951        deltas.push(OrderBookDelta::clear(instrument_id, 0, ts, ts));
1952        for (price, size) in bids {
1953            deltas.push(OrderBookDelta::new(
1954                instrument_id,
1955                BookAction::Add,
1956                BookOrder::new(
1957                    OrderSide::Buy,
1958                    Price::new(*price, 2),
1959                    Quantity::new(*size, 8),
1960                    0,
1961                ),
1962                0,
1963                0,
1964                ts,
1965                ts,
1966            ));
1967        }
1968
1969        for (price, size) in asks {
1970            deltas.push(OrderBookDelta::new(
1971                instrument_id,
1972                BookAction::Add,
1973                BookOrder::new(
1974                    OrderSide::Sell,
1975                    Price::new(*price, 2),
1976                    Quantity::new(*size, 8),
1977                    0,
1978                ),
1979                0,
1980                0,
1981                ts,
1982                ts,
1983            ));
1984        }
1985
1986        if let Some(last) = deltas.last_mut() {
1987            last.flags = RecordFlag::F_LAST as u8;
1988        }
1989
1990        book.apply_deltas(&OrderBookDeltas::new(instrument_id, deltas))
1991            .expect("failed to apply seed deltas");
1992        book
1993    }
1994
1995    fn crossing_bid_deltas(
1996        instrument_id: InstrumentId,
1997        bid_price: f64,
1998        bid_size: f64,
1999    ) -> OrderBookDeltas {
2000        let ts = UnixNanos::default();
2001        let delta = OrderBookDelta::new(
2002            instrument_id,
2003            BookAction::Add,
2004            BookOrder::new(
2005                OrderSide::Buy,
2006                Price::new(bid_price, 2),
2007                Quantity::new(bid_size, 8),
2008                0,
2009            ),
2010            RecordFlag::F_LAST as u8,
2011            0,
2012            ts,
2013            ts,
2014        );
2015        OrderBookDeltas::new(instrument_id, vec![delta])
2016    }
2017
2018    #[rstest]
2019    fn test_resolve_crossed_order_book_preserves_decimal_precision() {
2020        // Book seeded uncrossed: bid at 99.00 / ask at 100.05 size=0.50000000.
2021        // Venue delta adds a crossing bid at 100.10 size=1.00000001.
2022        // The reducing side (Buy) must end up at size = 0.50000001 exactly --
2023        // f64 subtraction of 1.00000001 - 0.5 would round this to 0.50000000 at 8 dp.
2024        let instrument = test_instrument();
2025        let instrument_id = instrument.id();
2026        let mut book = seed_book_with_levels(
2027            instrument_id,
2028            &[(99.00, 1.00000000)],
2029            &[(100.05, 0.50000000)],
2030        );
2031
2032        let venue_deltas = crossing_bid_deltas(instrument_id, 100.10, 1.00000001);
2033
2034        let resolved =
2035            DydxDataClient::resolve_crossed_order_book(&mut book, &venue_deltas, &instrument)
2036                .expect("resolution should succeed");
2037
2038        // An Update on the Buy side at the crossing price must carry the exact
2039        // Decimal-subtracted remainder (0.50000001), not the f64-rounded 0.50000000.
2040        let update = resolved
2041            .deltas
2042            .iter()
2043            .find(|d| {
2044                d.action == BookAction::Update
2045                    && d.order.side == Some(OrderSide::Buy)
2046                    && d.order.price.as_decimal() == dec!(100.10)
2047            })
2048            .expect("expected a Buy Update delta from crossed-book resolution");
2049        assert_eq!(update.order.size.as_decimal(), dec!(0.50000001));
2050
2051        assert_eq!(
2052            resolved.deltas.last().unwrap().flags,
2053            RecordFlag::F_LAST as u8,
2054        );
2055
2056        if let (Some(bid), Some(ask)) = (book.best_bid_price(), book.best_ask_price()) {
2057            assert!(bid < ask, "book still crossed: bid={bid:?} ask={ask:?}");
2058        }
2059    }
2060
2061    fn crossing_snapshot_batch(
2062        instrument_id: InstrumentId,
2063        bid_price: f64,
2064        bid_size: f64,
2065    ) -> OrderBookDeltas {
2066        let ts = UnixNanos::default();
2067        let snapshot = RecordFlag::F_SNAPSHOT as u8;
2068        let last = RecordFlag::F_LAST as u8;
2069        let deltas = vec![OrderBookDelta::new(
2070            instrument_id,
2071            BookAction::Add,
2072            BookOrder::new(
2073                OrderSide::Buy,
2074                Price::new(bid_price, 2),
2075                Quantity::new(bid_size, 8),
2076                0,
2077            ),
2078            snapshot | last,
2079            0,
2080            ts,
2081            ts,
2082        )];
2083        OrderBookDeltas::new(instrument_id, deltas)
2084    }
2085
2086    #[rstest]
2087    fn test_resolve_crossed_order_book_preserves_snapshot_flags() {
2088        let instrument = test_instrument();
2089        let instrument_id = instrument.id();
2090        let mut book = seed_book_with_levels(
2091            instrument_id,
2092            &[(99.00, 1.00000000)],
2093            &[(100.05, 0.50000000)],
2094        );
2095
2096        let venue_deltas = crossing_snapshot_batch(instrument_id, 100.10, 1.00000001);
2097
2098        let resolved =
2099            DydxDataClient::resolve_crossed_order_book(&mut book, &venue_deltas, &instrument)
2100                .expect("resolution should succeed");
2101
2102        let snapshot = RecordFlag::F_SNAPSHOT as u8;
2103        let last = RecordFlag::F_LAST as u8;
2104
2105        for (idx, delta) in resolved.deltas.iter().enumerate() {
2106            assert!(
2107                delta.flags & snapshot != 0,
2108                "delta at index {idx} lost F_SNAPSHOT: flags={:#010b}",
2109                delta.flags,
2110            );
2111        }
2112        assert_eq!(
2113            resolved.deltas.last().unwrap().flags,
2114            snapshot | last,
2115            "snapshot terminator must be F_SNAPSHOT | F_LAST",
2116        );
2117    }
2118
2119    #[rstest]
2120    fn test_resolve_crossed_order_book_equal_sizes_removes_both_levels() {
2121        let instrument = test_instrument();
2122        let instrument_id = instrument.id();
2123        let mut book = seed_book_with_levels(
2124            instrument_id,
2125            &[(99.00, 1.00000000)],
2126            &[(100.05, 1.00000000)],
2127        );
2128
2129        let venue_deltas = crossing_bid_deltas(instrument_id, 100.10, 1.00000000);
2130
2131        let resolved =
2132            DydxDataClient::resolve_crossed_order_book(&mut book, &venue_deltas, &instrument)
2133                .expect("resolution should succeed");
2134
2135        let deletes_count = resolved
2136            .deltas
2137            .iter()
2138            .filter(|d| {
2139                d.action == BookAction::Delete
2140                    && (d.order.price.as_decimal() == dec!(100.10)
2141                        || d.order.price.as_decimal() == dec!(100.05))
2142            })
2143            .count();
2144        assert_eq!(deletes_count, 2);
2145    }
2146}