1use 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#[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 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 #[must_use]
179 pub fn venue(&self) -> Venue {
180 *DYDX_VENUE
181 }
182
183 #[must_use]
185 pub fn config(&self) -> &DydxDataClientConfig {
186 &self.config
187 }
188
189 #[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 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 #[must_use]
1135 pub fn get_instrument(&self, instrument_id: &InstrumentId) -> Option<InstrumentAny> {
1136 self.instrument_cache.get(instrument_id)
1137 }
1138
1139 #[must_use]
1141 pub fn get_instruments(&self) -> Vec<InstrumentAny> {
1142 self.instrument_cache.all_instruments()
1143 }
1144
1145 pub fn cache_instrument(&self, instrument: InstrumentAny) {
1147 self.instrument_cache.insert_instrument_only(instrument);
1148 }
1149
1150 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 #[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 #[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 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 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 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 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 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 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}