1use std::{
22 path::PathBuf,
23 str::FromStr,
24 sync::{
25 Arc,
26 atomic::{AtomicBool, Ordering},
27 },
28 time::Duration,
29};
30
31use ahash::AHashMap;
32use databento::{dbn, live::Subscription};
33use indexmap::IndexMap;
34use nautilus_common::{
35 clients::DataClient,
36 live::runner::get_data_event_sender,
37 messages::{
38 DataEvent, DataResponse,
39 data::{
40 BarsResponse, BookDeltasResponse, BookDepthResponse, InstrumentResponse,
41 InstrumentsResponse, QuotesResponse, RequestBars, RequestBookDeltas, RequestBookDepth,
42 RequestInstrument, RequestInstruments, RequestQuotes, RequestTrades,
43 SubscribeBookDeltas, SubscribeInstrument, SubscribeInstrumentStatus, SubscribeQuotes,
44 SubscribeTrades, TradesResponse, UnsubscribeBookDeltas, UnsubscribeInstrumentStatus,
45 UnsubscribeQuotes, UnsubscribeTrades,
46 },
47 },
48};
49use nautilus_core::{
50 AtomicMap, DurationNanos, Params, UnixNanos,
51 datetime::datetime_to_unix_nanos,
52 time::{AtomicTime, get_atomic_clock_realtime},
53};
54use nautilus_live::task::TaskGroup;
55use nautilus_model::{
56 data::{CustomData, Data},
57 enums::BarAggregation,
58 identifiers::{ClientId, InstrumentId, Symbol, Venue},
59 instruments::{Instrument, InstrumentAny},
60};
61use parking_lot::Mutex;
62use tokio_util::sync::CancellationToken;
63
64use crate::{
65 common::{Credential, DATABENTO_VENUE},
66 historical::{DatabentoHistoricalClient, RangeQueryParams},
67 live::{DatabentoFeedHandler, DatabentoMessage, HandlerCommand},
68 loader::DatabentoDataLoader,
69 symbology::instrument_id_to_symbol_string,
70 types::{Dataset, PublisherId},
71};
72
73const PRICE_PRECISION_PARAM: &str = "price_precision";
74const SCHEMA_PARAM: &str = "schema";
75const QUOTE_SCHEMAS: &[dbn::Schema] = &[
76 dbn::Schema::Mbp1,
77 dbn::Schema::Bbo1S,
78 dbn::Schema::Bbo1M,
79 dbn::Schema::Cmbp1,
80 dbn::Schema::Cbbo1S,
81 dbn::Schema::Cbbo1M,
82 dbn::Schema::Tbbo,
83 dbn::Schema::Tcbbo,
84];
85const TRADE_SCHEMAS: &[dbn::Schema] = &[
86 dbn::Schema::Trades,
87 dbn::Schema::Tbbo,
88 dbn::Schema::Tcbbo,
89 dbn::Schema::Mbp1,
90 dbn::Schema::Cmbp1,
91];
92
93#[derive(Debug, Clone)]
95#[cfg_attr(
96 feature = "python",
97 pyo3::pyclass(module = "nautilus_trader.adapters.databento", from_py_object)
98)]
99#[cfg_attr(
100 feature = "python",
101 pyo3_stub_gen::derive::gen_stub_pyclass(module = "nautilus_trader.adapters.databento")
102)]
103pub struct DatabentoDataClientConfig {
104 pub(crate) credential: Credential,
106 pub publishers_filepath: PathBuf,
108 pub venue_dataset_map: IndexMap<String, String>,
110 pub use_exchange_as_venue: bool,
112 pub bars_timestamp_on_close: bool,
114 pub reconnect_timeout_mins: Option<u64>,
116}
117
118#[cfg(feature = "python")]
119nautilus_core::impl_pyo3_config_getters!(DatabentoDataClientConfig {
120 publishers_filepath: PathBuf,
121 use_exchange_as_venue: bool,
122 bars_timestamp_on_close: bool,
123 venue_dataset_map: IndexMap<String, String>,
124});
125
126impl DatabentoDataClientConfig {
127 #[must_use]
129 pub fn new(
130 api_key: impl Into<String>,
131 publishers_filepath: PathBuf,
132 use_exchange_as_venue: bool,
133 bars_timestamp_on_close: bool,
134 ) -> Self {
135 Self {
136 credential: Credential::new(api_key),
137 publishers_filepath,
138 venue_dataset_map: IndexMap::new(),
139 use_exchange_as_venue,
140 bars_timestamp_on_close,
141 reconnect_timeout_mins: Some(10), }
143 }
144
145 #[must_use]
147 pub fn api_key(&self) -> &str {
148 self.credential.api_key()
149 }
150
151 #[must_use]
153 pub fn api_key_masked(&self) -> String {
154 self.credential.api_key_masked()
155 }
156}
157
158#[cfg_attr(feature = "python", pyo3::pyclass)]
164#[cfg_attr(
165 feature = "python",
166 pyo3_stub_gen::derive::gen_stub_pyclass(module = "nautilus_trader.adapters.databento")
167)]
168#[derive(Debug)]
169pub struct DatabentoDataClient {
170 client_id: ClientId,
171 config: DatabentoDataClientConfig,
172 is_connected: AtomicBool,
173 historical: DatabentoHistoricalClient,
174 loader: DatabentoDataLoader,
175 cmd_channels: Arc<Mutex<AHashMap<String, tokio::sync::mpsc::UnboundedSender<HandlerCommand>>>>,
176 task_handles: TaskGroup,
177 cancellation_token: CancellationToken,
178 publisher_venue_map: Arc<IndexMap<PublisherId, Venue>>,
179 symbol_venue_map: Arc<AtomicMap<Symbol, Venue>>,
180 data_sender: tokio::sync::mpsc::UnboundedSender<DataEvent>,
181}
182
183impl DatabentoDataClient {
184 pub fn new(
190 client_id: ClientId,
191 config: DatabentoDataClientConfig,
192 clock: &'static AtomicTime,
193 ) -> anyhow::Result<Self> {
194 let historical = DatabentoHistoricalClient::new(
195 config.credential.clone(),
196 config.publishers_filepath.clone(),
197 clock,
198 config.use_exchange_as_venue,
199 )?;
200
201 let mut loader = DatabentoDataLoader::new(Some(config.publishers_filepath.clone()))?;
203 for (venue, dataset) in &config.venue_dataset_map {
204 loader.set_dataset_for_venue(
205 Dataset::from(dataset.as_str()),
206 Venue::from(venue.as_str()),
207 );
208 }
209
210 let file_content = std::fs::read_to_string(&config.publishers_filepath)?;
212 let publishers_vec: Vec<crate::types::DatabentoPublisher> =
213 serde_json::from_str(&file_content)?;
214
215 let publisher_venue_map = publishers_vec
216 .into_iter()
217 .map(|p| (p.publisher_id, Venue::from(p.venue.as_str())))
218 .collect::<IndexMap<u16, Venue>>();
219
220 let data_sender = get_data_event_sender();
221
222 let task_handles = TaskGroup::new();
223
224 Ok(Self {
225 client_id,
226 config,
227 is_connected: AtomicBool::new(false),
228 historical,
229 loader,
230 cmd_channels: Arc::new(Mutex::new(AHashMap::new())),
231 cancellation_token: task_handles.cancellation_token(),
232 task_handles,
233 publisher_venue_map: Arc::new(publisher_venue_map),
234 symbol_venue_map: Arc::new(AtomicMap::new()),
235 data_sender,
236 })
237 }
238
239 #[must_use]
241 pub fn api_key(&self) -> &str {
242 self.config.api_key()
243 }
244
245 #[must_use]
247 pub fn api_key_masked(&self) -> String {
248 self.config.api_key_masked()
249 }
250
251 fn get_dataset_for_venue(&self, venue: Venue) -> anyhow::Result<String> {
257 self.loader
258 .get_dataset_for_venue(&venue)
259 .map(ToString::to_string)
260 .ok_or_else(|| anyhow::anyhow!("No dataset found for venue: {venue}"))
261 }
262
263 fn get_or_create_feed_handler(&self, dataset: &str) -> bool {
265 let mut channels = self.cmd_channels.lock();
266
267 if !channels.contains_key(dataset) {
268 log::debug!("Creating new feed handler for dataset: {dataset}");
269 let cmd_tx = self.initialize_live_feed(dataset.to_string());
270 channels.insert(dataset.to_string(), cmd_tx);
271
272 log::debug!("Feed handler created for dataset: {dataset}, channel stored");
273 return true;
274 }
275
276 false
277 }
278
279 fn send_subscription_to_dataset(
280 &self,
281 dataset: &str,
282 price_precision: Option<(Symbol, u8)>,
283 subscription: Subscription,
284 start_after_subscribe: bool,
285 ) -> anyhow::Result<()> {
286 let tx = {
287 let channels = self.cmd_channels.lock();
288 channels
289 .get(dataset)
290 .cloned()
291 .ok_or_else(|| anyhow::anyhow!("No feed handler found for dataset: {dataset}"))?
292 };
293
294 send_subscription_commands(
295 &tx,
296 dataset,
297 price_precision,
298 subscription,
299 start_after_subscribe,
300 )
301 }
302
303 fn send_close_to_active_feeds(&self) {
304 let channels = self.cmd_channels.lock();
305 for (dataset, tx) in channels.iter() {
306 if let Err(e) = tx.send(HandlerCommand::Close) {
307 log::warn!("Failed to send close command to dataset {dataset}: {e}");
308 }
309 }
310 }
311
312 fn clear_feed_channels(&self) {
313 let mut channels = self.cmd_channels.lock();
314 channels.clear();
315 }
316
317 fn abort_active_tasks(&self) {
318 self.task_handles.begin_shutdown();
319 }
320
321 fn spawn_task<F>(&self, future: F)
322 where
323 F: std::future::Future<Output = ()> + Send + 'static,
324 {
325 if let Err(e) = self.task_handles.spawn(future) {
326 log::debug!("Skipping Databento task after shutdown began: {e}");
327 }
328 }
329
330 fn initialize_live_feed(
332 &self,
333 dataset: String,
334 ) -> tokio::sync::mpsc::UnboundedSender<HandlerCommand> {
335 let (cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
336 let (msg_tx, msg_rx) = tokio::sync::mpsc::unbounded_channel();
337 let feed_dataset = dataset.clone();
338 let feed_channels = self.cmd_channels.clone();
339
340 let mut feed_handler = DatabentoFeedHandler::new(
341 self.config.credential.clone(),
342 dataset,
343 cmd_rx,
344 msg_tx,
345 (*self.publisher_venue_map).clone(),
346 self.symbol_venue_map.clone(),
347 self.config.use_exchange_as_venue,
348 self.config.bars_timestamp_on_close,
349 self.config.reconnect_timeout_mins,
350 );
351
352 let feed_future = async move {
353 if let Err(e) = feed_handler.run().await {
354 log::error!("Feed handler error: {e}");
355 }
356 feed_channels.lock().remove(&feed_dataset);
357 };
358
359 let cancellation_token = self.cancellation_token.clone();
360 let data_sender = self.data_sender.clone();
361
362 let msg_future = async move {
364 let mut msg_rx = msg_rx;
365
366 loop {
367 tokio::select! {
368 msg = msg_rx.recv() => {
369 match msg {
370 Some(DatabentoMessage::Data(data)) => {
371 log::debug!("Received data: {data:?}");
372 if let Err(e) = data_sender.send(DataEvent::Data(data)) {
373 log::error!("Failed to send data event: {e}");
374 }
375 }
376 Some(DatabentoMessage::Instrument(instrument)) => {
377 log::debug!("Received instrument definition: {}", instrument.id());
378 if let Err(e) = data_sender.send(DataEvent::Instrument(*instrument)) {
379 log::error!("Failed to send instrument: {e}");
380 }
381 }
382 Some(DatabentoMessage::Status(status)) => {
383 log::debug!("Received status: {status:?}");
384 if let Err(e) =
385 data_sender.send(DataEvent::Data(Data::InstrumentStatus(status)))
386 {
387 log::error!("Failed to send status data event: {e}");
388 }
389 }
390 Some(DatabentoMessage::Imbalance(imbalance)) => {
391 log::debug!("Received imbalance: {imbalance:?}");
392 let data = Data::Custom(CustomData::from_arc(Arc::new(imbalance)));
393 if let Err(e) = data_sender.send(DataEvent::Data(data)) {
394 log::error!("Failed to send imbalance data event: {e}");
395 }
396 }
397 Some(DatabentoMessage::Statistics(statistics)) => {
398 log::debug!("Received statistics: {statistics:?}");
399 let data = Data::Custom(CustomData::from_arc(Arc::new(statistics)));
400 if let Err(e) = data_sender.send(DataEvent::Data(data)) {
401 log::error!("Failed to send statistics data event: {e}");
402 }
403 }
404 Some(DatabentoMessage::SubscriptionAck(ack)) => {
405 log::debug!("Received subscription ack: {}", ack.message);
406 }
407 Some(DatabentoMessage::Error(error)) => {
408 log::error!("Feed handler error: {error}");
409 }
410 Some(DatabentoMessage::Close) => {
411 log::debug!("Feed handler closed");
412 break;
413 }
414 None => {
415 log::debug!("Message channel closed");
416 break;
417 }
418 }
419 }
420 () = cancellation_token.cancelled() => {
421 log::debug!("Message processing cancelled");
422 break;
423 }
424 }
425 }
426 };
427
428 if let Err(e) = self.task_handles.spawn(feed_future) {
429 log::warn!("Skipping Databento feed task after shutdown began: {e}");
430 }
431
432 if let Err(e) = self.task_handles.spawn(msg_future) {
433 log::warn!("Skipping Databento message task after shutdown began: {e}");
434 }
435
436 cmd_tx
437 }
438}
439
440#[async_trait::async_trait(?Send)]
441impl DataClient for DatabentoDataClient {
442 fn client_id(&self) -> ClientId {
444 self.client_id
445 }
446
447 fn venue(&self) -> Option<Venue> {
449 None
450 }
451
452 fn start(&mut self) -> anyhow::Result<()> {
458 log::debug!("Starting");
459 Ok(())
460 }
461
462 fn stop(&mut self) -> anyhow::Result<()> {
468 log::debug!("Stopping");
469
470 self.send_close_to_active_feeds();
471 self.clear_feed_channels();
472 self.cancellation_token.cancel();
473 self.abort_active_tasks();
474 self.is_connected.store(false, Ordering::Relaxed);
475
476 Ok(())
477 }
478
479 fn reset(&mut self) -> anyhow::Result<()> {
480 log::debug!("Resetting");
481 self.send_close_to_active_feeds();
482 self.clear_feed_channels();
483 self.abort_active_tasks();
484 self.is_connected.store(false, Ordering::Relaxed);
485 Ok(())
486 }
487
488 fn dispose(&mut self) -> anyhow::Result<()> {
489 log::debug!("Disposing");
490 self.stop()
491 }
492
493 async fn connect(&mut self) -> anyhow::Result<()> {
494 log::debug!("Connecting...");
495
496 if !self.task_handles.is_open() {
497 self.task_handles
498 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
499 .await
500 .map_err(|e| anyhow::anyhow!("Failed to terminate Databento tasks: {e}"))?;
501 self.task_handles
502 .start_generation()
503 .map_err(|e| anyhow::anyhow!("Failed to start Databento task generation: {e}"))?;
504 self.cancellation_token = self.task_handles.cancellation_token();
505 }
506
507 self.is_connected.store(true, Ordering::Relaxed);
508
509 log::info!("Connected");
510 Ok(())
511 }
512
513 async fn disconnect(&mut self) -> anyhow::Result<()> {
514 log::debug!("Disconnecting...");
515
516 self.send_close_to_active_feeds();
517 self.clear_feed_channels();
518 self.task_handles.begin_shutdown();
519
520 let tasks_result = self
521 .task_handles
522 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
523 .await
524 .map_err(|e| anyhow::anyhow!("Failed to terminate Databento tasks: {e}"));
525
526 self.is_connected.store(false, Ordering::Relaxed);
527
528 log::info!("Disconnected");
529 tasks_result
530 }
531
532 fn is_connected(&self) -> bool {
534 self.is_connected.load(Ordering::Relaxed)
535 }
536
537 fn is_disconnected(&self) -> bool {
538 !self.is_connected()
539 }
540
541 fn subscribe_instrument(&mut self, cmd: SubscribeInstrument) -> anyhow::Result<()> {
547 let dataset = self.get_dataset_for_venue(cmd.instrument_id.venue)?;
548 let start_after_subscribe = self.get_or_create_feed_handler(&dataset);
549
550 self.symbol_venue_map
551 .insert(cmd.instrument_id.symbol, cmd.instrument_id.venue);
552 let symbol = cmd.instrument_id.symbol.to_string();
553
554 let subscription = Subscription::builder()
555 .schema(databento::dbn::Schema::Definition)
556 .symbols(symbol)
557 .build();
558
559 self.send_subscription_to_dataset(&dataset, None, subscription, start_after_subscribe)?;
560
561 Ok(())
562 }
563
564 fn subscribe_quotes(&mut self, cmd: SubscribeQuotes) -> anyhow::Result<()> {
570 let dataset = self.get_dataset_for_venue(cmd.instrument_id.venue)?;
571 let symbol = cmd.instrument_id.symbol.to_string();
572 let price_precision = price_precision_from_params(cmd.params.as_ref())?
573 .map(|precision| (cmd.instrument_id.symbol, precision));
574 let schema = schema_from_params(cmd.params.as_ref(), dbn::Schema::Mbp1, QUOTE_SCHEMAS)?;
575
576 let subscription = Subscription::builder()
577 .schema(schema)
578 .symbols(symbol)
579 .build();
580
581 let start_after_subscribe = self.get_or_create_feed_handler(&dataset);
582 self.symbol_venue_map
583 .insert(cmd.instrument_id.symbol, cmd.instrument_id.venue);
584
585 self.send_subscription_to_dataset(
586 &dataset,
587 price_precision,
588 subscription,
589 start_after_subscribe,
590 )?;
591
592 Ok(())
593 }
594
595 fn subscribe_trades(&mut self, cmd: SubscribeTrades) -> anyhow::Result<()> {
601 let dataset = self.get_dataset_for_venue(cmd.instrument_id.venue)?;
602 let symbol = cmd.instrument_id.symbol.to_string();
603 let price_precision = price_precision_from_params(cmd.params.as_ref())?
604 .map(|precision| (cmd.instrument_id.symbol, precision));
605 let schema = schema_from_params(cmd.params.as_ref(), dbn::Schema::Trades, TRADE_SCHEMAS)?;
606
607 let subscription = Subscription::builder()
608 .schema(schema)
609 .symbols(symbol)
610 .build();
611
612 let start_after_subscribe = self.get_or_create_feed_handler(&dataset);
613 self.symbol_venue_map
614 .insert(cmd.instrument_id.symbol, cmd.instrument_id.venue);
615
616 self.send_subscription_to_dataset(
617 &dataset,
618 price_precision,
619 subscription,
620 start_after_subscribe,
621 )?;
622
623 Ok(())
624 }
625
626 fn subscribe_book_deltas(&mut self, cmd: SubscribeBookDeltas) -> anyhow::Result<()> {
632 let dataset = self.get_dataset_for_venue(cmd.instrument_id.venue)?;
633 let start_after_subscribe = self.get_or_create_feed_handler(&dataset);
634
635 self.symbol_venue_map
636 .insert(cmd.instrument_id.symbol, cmd.instrument_id.venue);
637 let symbol = cmd.instrument_id.symbol.to_string();
638
639 let subscription = Subscription::builder()
640 .schema(databento::dbn::Schema::Mbo) .symbols(symbol)
642 .build();
643
644 self.send_subscription_to_dataset(&dataset, None, subscription, start_after_subscribe)?;
645
646 Ok(())
647 }
648
649 fn subscribe_instrument_status(
655 &mut self,
656 cmd: SubscribeInstrumentStatus,
657 ) -> anyhow::Result<()> {
658 let dataset = self.get_dataset_for_venue(cmd.instrument_id.venue)?;
659 let start_after_subscribe = self.get_or_create_feed_handler(&dataset);
660
661 self.symbol_venue_map
662 .insert(cmd.instrument_id.symbol, cmd.instrument_id.venue);
663 let symbol = cmd.instrument_id.symbol.to_string();
664
665 let subscription = Subscription::builder()
666 .schema(databento::dbn::Schema::Status)
667 .symbols(symbol)
668 .build();
669
670 self.send_subscription_to_dataset(&dataset, None, subscription, start_after_subscribe)?;
671
672 Ok(())
673 }
674
675 fn unsubscribe_quotes(&mut self, cmd: &UnsubscribeQuotes) -> anyhow::Result<()> {
677 log::warn!(
681 "Databento does not support granular unsubscribing - ignoring unsubscribe request for {}",
682 cmd.instrument_id
683 );
684
685 Ok(())
686 }
687
688 fn unsubscribe_trades(&mut self, cmd: &UnsubscribeTrades) -> anyhow::Result<()> {
689 log::warn!(
693 "Databento does not support granular unsubscribing - ignoring unsubscribe request for {}",
694 cmd.instrument_id
695 );
696
697 Ok(())
698 }
699
700 fn unsubscribe_book_deltas(&mut self, cmd: &UnsubscribeBookDeltas) -> anyhow::Result<()> {
701 log::warn!(
705 "Databento does not support granular unsubscribing - ignoring unsubscribe request for {}",
706 cmd.instrument_id
707 );
708
709 Ok(())
710 }
711
712 fn unsubscribe_instrument_status(
713 &mut self,
714 cmd: &UnsubscribeInstrumentStatus,
715 ) -> anyhow::Result<()> {
716 log::warn!(
720 "Databento does not support granular unsubscribing - ignoring unsubscribe request for {}",
721 cmd.instrument_id
722 );
723
724 Ok(())
725 }
726
727 fn request_instruments(&self, request: RequestInstruments) -> anyhow::Result<()> {
728 log::debug!("Request instruments: {request:?}");
729
730 let historical_client = self.historical.clone();
731 let data_sender = self.data_sender.clone();
732 let dataset = request
733 .venue
734 .map(|venue| self.get_dataset_for_venue(venue))
735 .transpose()?
736 .unwrap_or_else(|| "GLBX.MDP3".to_string());
737 let request_id = request.request_id;
738 let client_id = request.client_id.unwrap_or(self.client_id);
739 let venue = request.venue.unwrap_or(*DATABENTO_VENUE);
740 let start_nanos = datetime_to_unix_nanos(request.start);
741 let end_nanos = datetime_to_unix_nanos(request.end);
742 let request_params = request.params;
743 let (query_start, query_end) = resolve_request_time_range(start_nanos, end_nanos);
744
745 self.spawn_task(async move {
746 let query_params = instruments_query_params(dataset, query_start, query_end);
747
748 match historical_client.get_range_instruments(query_params).await {
749 Ok(instruments) => {
750 log::debug!("Retrieved {} instruments", instruments.len());
751
752 let response = DataResponse::Instruments(InstrumentsResponse::new(
753 request_id,
754 client_id,
755 venue,
756 instruments,
757 start_nanos,
758 end_nanos,
759 get_atomic_clock_realtime().get_time_ns(),
760 request_params,
761 ));
762
763 if let Err(e) = data_sender.send(DataEvent::Response(response)) {
764 log::error!("Failed to send instruments response: {e}");
765 }
766 }
767 Err(e) => {
768 log::error!("Failed to request instruments: {e}");
769 let response = DataResponse::Instruments(InstrumentsResponse::new(
770 request_id,
771 client_id,
772 venue,
773 Vec::new(),
774 start_nanos,
775 end_nanos,
776 get_atomic_clock_realtime().get_time_ns(),
777 request_params,
778 ));
779
780 send_data_response(&data_sender, response, "empty instruments");
781 }
782 }
783 });
784
785 Ok(())
786 }
787
788 fn request_instrument(&self, request: RequestInstrument) -> anyhow::Result<()> {
789 log::debug!("Request instrument: {request:?}");
790
791 let dataset = self.get_dataset_for_venue(request.instrument_id.venue)?;
792 let historical_client = self.historical.clone();
793 let data_sender = self.data_sender.clone();
794 let instrument_id = request.instrument_id;
795 let request_id = request.request_id;
796 let client_id = request.client_id.unwrap_or(self.client_id);
797 let start_nanos = datetime_to_unix_nanos(request.start);
798 let end_nanos = datetime_to_unix_nanos(request.end);
799 let request_params = request.params;
800 let (query_start, query_end) = resolve_request_time_range(start_nanos, end_nanos);
801
802 self.spawn_task(async move {
803 let query_params =
804 instrument_query_params(dataset, instrument_id, query_start, query_end);
805
806 match historical_client.get_range_instruments(query_params).await {
807 Ok(instruments) => {
808 let instrument = requested_instrument(instruments, instrument_id);
809
810 let Some(instrument) = instrument else {
811 log::error!("Instrument not found: {instrument_id}");
812 return;
813 };
814
815 let response = DataResponse::Instrument(Box::new(InstrumentResponse::new(
816 request_id,
817 client_id,
818 instrument.id(),
819 instrument,
820 start_nanos,
821 end_nanos,
822 get_atomic_clock_realtime().get_time_ns(),
823 request_params,
824 )));
825
826 if let Err(e) = data_sender.send(DataEvent::Response(response)) {
827 log::error!("Failed to send instrument response: {e}");
828 }
829 }
830 Err(e) => {
831 log::error!("Failed to request instrument {instrument_id}: {e}");
832 }
833 }
834 });
835
836 Ok(())
837 }
838
839 fn request_quotes(&self, request: RequestQuotes) -> anyhow::Result<()> {
840 log::debug!("Request quotes: {request:?}");
841
842 let historical_client = self.historical.clone();
843 let data_sender = self.data_sender.clone();
844 let dataset = self.get_dataset_for_venue(request.instrument_id.venue)?;
845 let instrument_id = request.instrument_id;
846 let symbols = historical_client.prepare_symbols_from_instrument_ids(&[instrument_id]);
847 let request_id = request.request_id;
848 let client_id = request.client_id.unwrap_or(self.client_id);
849 let start_nanos = datetime_to_unix_nanos(request.start);
850 let end_nanos = datetime_to_unix_nanos(request.end);
851 let limit = request.limit.map(|limit| limit.get() as u64);
852 let request_params = request.params;
853 let price_precision = price_precision_from_params(request_params.as_ref())?;
854 let schema = schema_from_params(request_params.as_ref(), dbn::Schema::Mbp1, QUOTE_SCHEMAS)?
855 .to_string();
856 let (query_start, query_end) = resolve_request_time_range(start_nanos, end_nanos);
857
858 self.spawn_task(async move {
859 seed_price_precision_if_needed(
860 &historical_client,
861 dataset.as_str(),
862 instrument_id,
863 query_start,
864 query_end,
865 price_precision,
866 )
867 .await;
868
869 let params = RangeQueryParams {
870 dataset,
871 symbols,
872 start: query_start,
873 end: query_end,
874 limit,
875 price_precision,
876 };
877
878 match historical_client
879 .get_range_quotes(params, Some(schema))
880 .await
881 {
882 Ok(quotes) => {
883 log::debug!("Retrieved {} quotes", quotes.len());
884 let response = DataResponse::Quotes(QuotesResponse::new(
885 request_id,
886 client_id,
887 instrument_id,
888 quotes,
889 start_nanos,
890 end_nanos,
891 get_atomic_clock_realtime().get_time_ns(),
892 request_params,
893 ));
894
895 if let Err(e) = data_sender.send(DataEvent::Response(response)) {
896 log::error!("Failed to send quotes response: {e}");
897 }
898 }
899 Err(e) => {
900 log::error!("Failed to request quotes: {e}");
901 let response = DataResponse::Quotes(QuotesResponse::new(
902 request_id,
903 client_id,
904 instrument_id,
905 Vec::new(),
906 start_nanos,
907 end_nanos,
908 get_atomic_clock_realtime().get_time_ns(),
909 request_params,
910 ));
911
912 send_data_response(&data_sender, response, "empty quotes");
913 }
914 }
915 });
916
917 Ok(())
918 }
919
920 fn request_trades(&self, request: RequestTrades) -> anyhow::Result<()> {
921 log::debug!("Request trades: {request:?}");
922
923 let historical_client = self.historical.clone();
924 let data_sender = self.data_sender.clone();
925 let dataset = self.get_dataset_for_venue(request.instrument_id.venue)?;
926 let instrument_id = request.instrument_id;
927 let symbols = historical_client.prepare_symbols_from_instrument_ids(&[instrument_id]);
928 let request_id = request.request_id;
929 let client_id = request.client_id.unwrap_or(self.client_id);
930 let start_nanos = datetime_to_unix_nanos(request.start);
931 let end_nanos = datetime_to_unix_nanos(request.end);
932 let limit = request.limit.map(|limit| limit.get() as u64);
933 let request_params = request.params;
934 let price_precision = price_precision_from_params(request_params.as_ref())?;
935 let schema =
936 schema_from_params(request_params.as_ref(), dbn::Schema::Trades, TRADE_SCHEMAS)?
937 .to_string();
938 let (query_start, query_end) = resolve_request_time_range(start_nanos, end_nanos);
939
940 self.spawn_task(async move {
941 seed_price_precision_if_needed(
942 &historical_client,
943 dataset.as_str(),
944 instrument_id,
945 query_start,
946 query_end,
947 price_precision,
948 )
949 .await;
950
951 let params = RangeQueryParams {
952 dataset,
953 symbols,
954 start: query_start,
955 end: query_end,
956 limit,
957 price_precision,
958 };
959
960 match historical_client
961 .get_range_trades(params, Some(schema))
962 .await
963 {
964 Ok(trades) => {
965 log::debug!("Retrieved {} trades", trades.len());
966 let response = DataResponse::Trades(TradesResponse::new(
967 request_id,
968 client_id,
969 instrument_id,
970 trades,
971 start_nanos,
972 end_nanos,
973 get_atomic_clock_realtime().get_time_ns(),
974 request_params,
975 ));
976
977 if let Err(e) = data_sender.send(DataEvent::Response(response)) {
978 log::error!("Failed to send trades response: {e}");
979 }
980 }
981 Err(e) => {
982 log::error!("Failed to request trades: {e}");
983 let response = DataResponse::Trades(TradesResponse::new(
984 request_id,
985 client_id,
986 instrument_id,
987 Vec::new(),
988 start_nanos,
989 end_nanos,
990 get_atomic_clock_realtime().get_time_ns(),
991 request_params,
992 ));
993
994 send_data_response(&data_sender, response, "empty trades");
995 }
996 }
997 });
998
999 Ok(())
1000 }
1001
1002 fn request_bars(&self, request: RequestBars) -> anyhow::Result<()> {
1003 log::debug!("Request bars: {request:?}");
1004
1005 let historical_client = self.historical.clone();
1006 let data_sender = self.data_sender.clone();
1007 let instrument_id = request.bar_type.instrument_id();
1008 let dataset = self.get_dataset_for_venue(instrument_id.venue)?;
1009 let symbols = historical_client.prepare_symbols_from_instrument_ids(&[instrument_id]);
1010 let request_id = request.request_id;
1011 let client_id = request.client_id.unwrap_or(self.client_id);
1012 let bar_type = request.bar_type;
1013 let start_nanos = datetime_to_unix_nanos(request.start);
1014 let end_nanos = datetime_to_unix_nanos(request.end);
1015 let limit = request.limit.map(|limit| limit.get() as u64);
1016 let request_params = request.params;
1017 let price_precision = price_precision_from_params(request_params.as_ref())?;
1018 let timestamp_on_close = self.config.bars_timestamp_on_close;
1019 let (query_start, query_end) = resolve_request_time_range(start_nanos, end_nanos);
1020
1021 self.spawn_task(async move {
1022 seed_price_precision_if_needed(
1023 &historical_client,
1024 dataset.as_str(),
1025 instrument_id,
1026 query_start,
1027 query_end,
1028 price_precision,
1029 )
1030 .await;
1031
1032 let params = RangeQueryParams {
1033 dataset,
1034 symbols,
1035 start: query_start,
1036 end: query_end,
1037 limit,
1038 price_precision,
1039 };
1040
1041 let aggregation = match bar_type.spec().aggregation {
1042 BarAggregation::Second => BarAggregation::Second,
1043 BarAggregation::Minute => BarAggregation::Minute,
1044 BarAggregation::Hour => BarAggregation::Hour,
1045 BarAggregation::Day => BarAggregation::Day,
1046 _ => {
1047 log::error!(
1048 "Unsupported bar aggregation: {:?}",
1049 bar_type.spec().aggregation
1050 );
1051 let response = DataResponse::Bars(BarsResponse::new(
1052 request_id,
1053 client_id,
1054 bar_type,
1055 Vec::new(),
1056 start_nanos,
1057 end_nanos,
1058 get_atomic_clock_realtime().get_time_ns(),
1059 request_params,
1060 ));
1061
1062 send_data_response(&data_sender, response, "empty bars");
1063 return;
1064 }
1065 };
1066
1067 match historical_client
1068 .get_range_bars(params, aggregation, timestamp_on_close)
1069 .await
1070 {
1071 Ok(bars) => {
1072 log::debug!("Retrieved {} bars", bars.len());
1073 let response = DataResponse::Bars(BarsResponse::new(
1074 request_id,
1075 client_id,
1076 bar_type,
1077 bars,
1078 start_nanos,
1079 end_nanos,
1080 get_atomic_clock_realtime().get_time_ns(),
1081 request_params,
1082 ));
1083
1084 if let Err(e) = data_sender.send(DataEvent::Response(response)) {
1085 log::error!("Failed to send bars response: {e}");
1086 }
1087 }
1088 Err(e) => {
1089 log::error!("Failed to request bars: {e}");
1090 let response = DataResponse::Bars(BarsResponse::new(
1091 request_id,
1092 client_id,
1093 bar_type,
1094 Vec::new(),
1095 start_nanos,
1096 end_nanos,
1097 get_atomic_clock_realtime().get_time_ns(),
1098 request_params,
1099 ));
1100
1101 send_data_response(&data_sender, response, "empty bars");
1102 }
1103 }
1104 });
1105
1106 Ok(())
1107 }
1108
1109 fn request_book_depth(&self, request: RequestBookDepth) -> anyhow::Result<()> {
1110 log::debug!("Request book depth: {request:?}");
1111
1112 let historical_client = self.historical.clone();
1113 let data_sender = self.data_sender.clone();
1114 let dataset = self.get_dataset_for_venue(request.instrument_id.venue)?;
1115 let instrument_id = request.instrument_id;
1116 let symbols = historical_client.prepare_symbols_from_instrument_ids(&[instrument_id]);
1117 let request_id = request.request_id;
1118 let client_id = request.client_id.unwrap_or(self.client_id);
1119 let start_nanos = datetime_to_unix_nanos(request.start);
1120 let end_nanos = datetime_to_unix_nanos(request.end);
1121 let limit = request.limit.map(|limit| limit.get() as u64);
1122 let depth = request.depth.map(|depth| depth.get());
1123 let request_params = request.params;
1124 let price_precision = price_precision_from_params(request_params.as_ref())?;
1125 let (query_start, query_end) = resolve_request_time_range(start_nanos, end_nanos);
1126
1127 self.spawn_task(async move {
1128 seed_price_precision_if_needed(
1129 &historical_client,
1130 dataset.as_str(),
1131 instrument_id,
1132 query_start,
1133 query_end,
1134 price_precision,
1135 )
1136 .await;
1137
1138 let params = RangeQueryParams {
1139 dataset,
1140 symbols,
1141 start: query_start,
1142 end: query_end,
1143 limit,
1144 price_precision,
1145 };
1146
1147 match historical_client
1148 .get_range_order_book_depth10(params, depth)
1149 .await
1150 {
1151 Ok(depths) => {
1152 log::debug!("Retrieved {} order book depths", depths.len());
1153 let response = DataResponse::BookDepth(BookDepthResponse::new(
1154 request_id,
1155 client_id,
1156 instrument_id,
1157 depths,
1158 start_nanos,
1159 end_nanos,
1160 get_atomic_clock_realtime().get_time_ns(),
1161 request_params,
1162 ));
1163
1164 send_data_response(&data_sender, response, "book depth");
1165 }
1166 Err(e) => {
1167 log::error!("Failed to request order book depths: {e}");
1168 let response = DataResponse::BookDepth(BookDepthResponse::new(
1169 request_id,
1170 client_id,
1171 instrument_id,
1172 Vec::new(),
1173 start_nanos,
1174 end_nanos,
1175 get_atomic_clock_realtime().get_time_ns(),
1176 request_params,
1177 ));
1178
1179 send_data_response(&data_sender, response, "empty book depth");
1180 }
1181 }
1182 });
1183
1184 Ok(())
1185 }
1186
1187 fn request_book_deltas(&self, request: RequestBookDeltas) -> anyhow::Result<()> {
1188 log::debug!("Request book deltas: {request:?}");
1189
1190 let historical_client = self.historical.clone();
1191 let data_sender = self.data_sender.clone();
1192 let dataset = self.get_dataset_for_venue(request.instrument_id.venue)?;
1193 let instrument_id = request.instrument_id;
1194 let symbols = historical_client.prepare_symbols_from_instrument_ids(&[instrument_id]);
1195 let request_id = request.request_id;
1196 let client_id = request.client_id.unwrap_or(self.client_id);
1197 let start_nanos = datetime_to_unix_nanos(request.start);
1198 let end_nanos = datetime_to_unix_nanos(request.end);
1199 let limit = request.limit.map(|limit| limit.get() as u64);
1200 let request_params = request.params;
1201 let price_precision = price_precision_from_params(request_params.as_ref())?;
1202 let (query_start, query_end) = resolve_request_time_range(start_nanos, end_nanos);
1203
1204 self.spawn_task(async move {
1205 seed_price_precision_if_needed(
1206 &historical_client,
1207 dataset.as_str(),
1208 instrument_id,
1209 query_start,
1210 query_end,
1211 price_precision,
1212 )
1213 .await;
1214
1215 let params = RangeQueryParams {
1216 dataset,
1217 symbols,
1218 start: query_start,
1219 end: query_end,
1220 limit,
1221 price_precision,
1222 };
1223
1224 match historical_client.get_range_order_book_deltas(params).await {
1225 Ok(deltas) => {
1226 log::debug!("Retrieved {} order book deltas", deltas.len());
1227 let response = BookDeltasResponse::new(
1228 request_id,
1229 client_id,
1230 instrument_id,
1231 deltas,
1232 start_nanos,
1233 end_nanos,
1234 get_atomic_clock_realtime().get_time_ns(),
1235 request_params,
1236 );
1237
1238 for response in partition_book_deltas_response(response) {
1239 send_data_response(
1240 &data_sender,
1241 DataResponse::BookDeltas(response),
1242 "book deltas",
1243 );
1244 }
1245 }
1246 Err(e) => {
1247 log::error!("Failed to request order book deltas: {e}");
1248 let response = DataResponse::BookDeltas(BookDeltasResponse::new(
1249 request_id,
1250 client_id,
1251 instrument_id,
1252 Vec::new(),
1253 start_nanos,
1254 end_nanos,
1255 get_atomic_clock_realtime().get_time_ns(),
1256 request_params,
1257 ));
1258
1259 send_data_response(&data_sender, response, "empty book deltas");
1260 }
1261 }
1262 });
1263
1264 Ok(())
1265 }
1266}
1267
1268fn instruments_query_params(
1269 dataset: String,
1270 start_nanos: UnixNanos,
1271 end_nanos: Option<UnixNanos>,
1272) -> RangeQueryParams {
1273 RangeQueryParams {
1274 dataset,
1275 symbols: vec!["ALL_SYMBOLS".to_string()],
1276 start: start_nanos,
1277 end: end_nanos,
1278 limit: None,
1279 price_precision: None,
1280 }
1281}
1282
1283fn instrument_query_params(
1284 dataset: String,
1285 instrument_id: InstrumentId,
1286 start_nanos: UnixNanos,
1287 end_nanos: Option<UnixNanos>,
1288) -> RangeQueryParams {
1289 RangeQueryParams {
1290 dataset,
1291 symbols: vec![instrument_id_to_symbol_string(
1292 instrument_id,
1293 &mut AHashMap::new(),
1294 )],
1295 start: start_nanos,
1296 end: end_nanos,
1297 limit: None,
1298 price_precision: None,
1299 }
1300}
1301
1302fn resolve_request_time_range(
1303 start_nanos: Option<UnixNanos>,
1304 end_nanos: Option<UnixNanos>,
1305) -> (UnixNanos, Option<UnixNanos>) {
1306 let mut end = end_nanos.unwrap_or_else(|| get_atomic_clock_realtime().get_time_ns());
1307 let mut start = start_nanos.unwrap_or_else(|| start_of_utc_day(end));
1308
1309 if start > end {
1310 start = end;
1311 }
1312
1313 if start == end {
1314 if end.is_zero() {
1315 end += DurationNanos::new(1);
1316 } else {
1317 start -= DurationNanos::new(1);
1318 }
1319 }
1320
1321 (start, Some(end))
1322}
1323
1324fn start_of_utc_day(timestamp: UnixNanos) -> UnixNanos {
1325 timestamp.floor(DurationNanos::from_days(1))
1326}
1327
1328async fn seed_price_precision_if_needed(
1329 historical_client: &DatabentoHistoricalClient,
1330 dataset: &str,
1331 instrument_id: InstrumentId,
1332 start_nanos: UnixNanos,
1333 end_nanos: Option<UnixNanos>,
1334 price_precision: Option<u8>,
1335) {
1336 if price_precision.is_some()
1337 || historical_client
1338 .price_precision(instrument_id.symbol)
1339 .is_some()
1340 {
1341 return;
1342 }
1343
1344 let query_params =
1345 instrument_query_params(dataset.to_string(), instrument_id, start_nanos, end_nanos);
1346
1347 if let Err(e) = historical_client.get_range_instruments(query_params).await {
1348 log::warn!("Failed to seed price precision for {instrument_id}: {e}");
1349 }
1350}
1351
1352fn send_data_response(
1353 data_sender: &tokio::sync::mpsc::UnboundedSender<DataEvent>,
1354 response: DataResponse,
1355 label: &str,
1356) {
1357 if let Err(e) = data_sender.send(DataEvent::Response(response)) {
1358 log::error!("Failed to send {label} response: {e}");
1359 }
1360}
1361
1362fn partition_book_deltas_response(mut response: BookDeltasResponse) -> Vec<BookDeltasResponse> {
1363 if response.data.is_empty() {
1364 return vec![response];
1365 }
1366
1367 let mut partitions = IndexMap::new();
1368
1369 for delta in std::mem::take(&mut response.data) {
1370 partitions
1371 .entry(delta.instrument_id)
1372 .or_insert_with(Vec::new)
1373 .push(delta);
1374 }
1375
1376 partitions
1377 .into_iter()
1378 .map(|(instrument_id, data)| {
1379 let mut child = response.clone();
1380 child.instrument_id = instrument_id;
1381 child.data = data;
1382 child
1383 })
1384 .collect()
1385}
1386
1387fn requested_instrument(
1388 instruments: Vec<InstrumentAny>,
1389 instrument_id: InstrumentId,
1390) -> Option<InstrumentAny> {
1391 instruments
1392 .into_iter()
1393 .rev()
1394 .find(|instrument| instrument.id() == instrument_id)
1395}
1396
1397fn price_precision_from_params(params: Option<&Params>) -> anyhow::Result<Option<u8>> {
1398 let Some(price_precision) = params.and_then(|params| params.get_u64(PRICE_PRECISION_PARAM))
1399 else {
1400 return Ok(None);
1401 };
1402
1403 Ok(Some(u8::try_from(price_precision).map_err(|_| {
1404 anyhow::anyhow!(
1405 "`{PRICE_PRECISION_PARAM}` must be less than or equal to {}",
1406 u8::MAX
1407 )
1408 })?))
1409}
1410
1411fn schema_from_params(
1412 params: Option<&Params>,
1413 default_schema: dbn::Schema,
1414 allowed_schemas: &[dbn::Schema],
1415) -> anyhow::Result<dbn::Schema> {
1416 let schema = if let Some(schema) = params.and_then(|params| params.get_str(SCHEMA_PARAM)) {
1417 dbn::Schema::from_str(schema)?
1418 } else {
1419 default_schema
1420 };
1421
1422 if allowed_schemas.contains(&schema) {
1423 return Ok(schema);
1424 }
1425
1426 let allowed = allowed_schemas
1427 .iter()
1428 .map(dbn::Schema::as_str)
1429 .collect::<Vec<_>>()
1430 .join(", ");
1431 anyhow::bail!(
1432 "Invalid `{SCHEMA_PARAM}` '{}'. Must be one of: {allowed}",
1433 schema.as_str()
1434 );
1435}
1436
1437fn send_subscription_commands(
1438 tx: &tokio::sync::mpsc::UnboundedSender<HandlerCommand>,
1439 dataset: &str,
1440 price_precision: Option<(Symbol, u8)>,
1441 subscription: Subscription,
1442 start_after_subscribe: bool,
1443) -> anyhow::Result<()> {
1444 if let Some((symbol, precision)) = price_precision {
1445 tx.send(HandlerCommand::SetPricePrecision(symbol, precision))
1446 .map_err(|e| anyhow::anyhow!("Failed to send command to dataset {dataset}: {e}"))?;
1447 }
1448
1449 tx.send(HandlerCommand::Subscribe(subscription))
1450 .map_err(|e| anyhow::anyhow!("Failed to send command to dataset {dataset}: {e}"))?;
1451
1452 if start_after_subscribe {
1453 tx.send(HandlerCommand::Start)
1454 .map_err(|e| anyhow::anyhow!("Failed to send command to dataset {dataset}: {e}"))?;
1455 }
1456
1457 Ok(())
1458}
1459
1460#[cfg(test)]
1461mod tests {
1462 use std::path::PathBuf;
1463
1464 use nautilus_common::live::runner::replace_data_event_sender;
1465 use nautilus_core::UUID4;
1466 use nautilus_model::{
1467 data::OrderBookDelta,
1468 identifiers::{ClientId, InstrumentId},
1469 instruments::{CurrencyPair, InstrumentAny},
1470 types::{Currency, Price, Quantity},
1471 };
1472 use rstest::rstest;
1473 use serde_json::json;
1474
1475 use super::*;
1476
1477 #[derive(Clone, Copy)]
1478 enum SubscribeKind {
1479 Quotes,
1480 Trades,
1481 }
1482
1483 fn currency_pair(instrument_id: &str) -> InstrumentAny {
1484 currency_pair_with_ts_init(instrument_id, UnixNanos::default())
1485 }
1486
1487 fn test_data_client() -> DatabentoDataClient {
1488 let (sender, _receiver) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
1489 replace_data_event_sender(sender);
1490
1491 let config = DatabentoDataClientConfig::new(
1492 "32-character-with-lots-of-filler",
1493 PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("publishers.json"),
1494 true,
1495 true,
1496 );
1497 DatabentoDataClient::new(
1498 ClientId::from("DATABENTO-TEST"),
1499 config,
1500 get_atomic_clock_realtime(),
1501 )
1502 .expect("test client should initialize")
1503 }
1504
1505 #[rstest]
1506 #[tokio::test]
1507 async fn test_stop_closes_task_admission_until_disconnect_drains() {
1508 let mut client = test_data_client();
1509
1510 client
1511 .task_handles
1512 .spawn(async { std::future::pending::<()>().await })
1513 .unwrap();
1514 client.is_connected.store(true, Ordering::Relaxed);
1515
1516 client.stop().unwrap();
1517
1518 assert!(!client.task_handles.is_open());
1519 assert!(client.is_disconnected());
1520
1521 client.disconnect().await.unwrap();
1522 assert!(client.task_handles.is_empty());
1523 }
1524
1525 #[rstest]
1526 #[case("EQUS", "EQUS.PLUS")] #[case("GLBX", "EQUS.MINI")] fn test_venue_dataset_map_overrides_default(#[case] venue: &str, #[case] dataset: &str) {
1529 let (sender, _receiver) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
1530 replace_data_event_sender(sender);
1531
1532 let mut config = DatabentoDataClientConfig::new(
1533 "32-character-with-lots-of-filler",
1534 PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("publishers.json"),
1535 true,
1536 true,
1537 );
1538 config.venue_dataset_map = IndexMap::from([(venue.to_string(), dataset.to_string())]);
1539
1540 let client = DatabentoDataClient::new(
1541 ClientId::from("DATABENTO-TEST"),
1542 config,
1543 get_atomic_clock_realtime(),
1544 )
1545 .expect("test client should initialize");
1546
1547 assert_eq!(
1548 client.get_dataset_for_venue(Venue::from(venue)).unwrap(),
1549 dataset
1550 );
1551
1552 assert_eq!(
1554 client.get_dataset_for_venue(Venue::from("XCBO")).unwrap(),
1555 "OPRA.PILLAR"
1556 );
1557 }
1558
1559 fn subscribe_quotes_cmd(params: Option<Params>) -> SubscribeQuotes {
1560 SubscribeQuotes::new(
1561 InstrumentId::from("ESM4.GLBX"),
1562 Some(ClientId::from("DATABENTO-TEST")),
1563 None,
1564 UUID4::new(),
1565 UnixNanos::default(),
1566 None,
1567 params,
1568 )
1569 }
1570
1571 fn subscribe_trades_cmd(params: Option<Params>) -> SubscribeTrades {
1572 SubscribeTrades::new(
1573 InstrumentId::from("ESM4.GLBX"),
1574 Some(ClientId::from("DATABENTO-TEST")),
1575 None,
1576 UUID4::new(),
1577 UnixNanos::default(),
1578 None,
1579 params,
1580 )
1581 }
1582
1583 fn currency_pair_with_ts_init(instrument_id: &str, ts_init: UnixNanos) -> InstrumentAny {
1584 let instrument_id = InstrumentId::from(instrument_id);
1585 InstrumentAny::CurrencyPair(
1586 CurrencyPair::builder()
1587 .instrument_id(instrument_id)
1588 .raw_symbol(instrument_id.symbol)
1589 .base_currency(Currency::from("BTC"))
1590 .quote_currency(Currency::from("USDT"))
1591 .price_precision(2)
1592 .size_precision(6)
1593 .price_increment(Price::from("0.01"))
1594 .size_increment(Quantity::from("0.000001"))
1595 .ts_event(UnixNanos::default())
1596 .ts_init(ts_init)
1597 .build()
1598 .unwrap(),
1599 )
1600 }
1601
1602 #[rstest]
1603 fn test_instruments_query_params_requests_all_symbols() {
1604 let start = UnixNanos::from(1_000_000_000);
1605 let end = UnixNanos::from(2_000_000_000);
1606
1607 let params = instruments_query_params("GLBX.MDP3".to_string(), start, Some(end));
1608
1609 assert_eq!(params.dataset, "GLBX.MDP3");
1610 assert_eq!(params.symbols, vec!["ALL_SYMBOLS"]);
1611 assert_eq!(params.start, start);
1612 assert_eq!(params.end, Some(end));
1613 assert_eq!(params.limit, None);
1614 assert_eq!(params.price_precision, None);
1615 }
1616
1617 #[rstest]
1618 fn test_instrument_query_params_requests_single_symbol() {
1619 let instrument_id = InstrumentId::from("ESM4.GLBX");
1620
1621 let start = UnixNanos::from(1_000_000_000);
1622 let end = UnixNanos::from(2_000_000_000);
1623
1624 let params =
1625 instrument_query_params("GLBX.MDP3".to_string(), instrument_id, start, Some(end));
1626
1627 assert_eq!(params.dataset, "GLBX.MDP3");
1628 assert_eq!(params.symbols, vec!["ESM4"]);
1629 assert_eq!(params.start, start);
1630 assert_eq!(params.end, Some(end));
1631 assert_eq!(params.limit, None);
1632 assert_eq!(params.price_precision, None);
1633 }
1634
1635 #[rstest]
1636 fn test_resolve_request_time_range_defaults_to_end_day() {
1637 let end = UnixNanos::from(1_706_443_200_000_000_001);
1638
1639 let (start, resolved_end) = resolve_request_time_range(None, Some(end));
1640
1641 assert_eq!(start, UnixNanos::from(1_706_400_000_000_000_000));
1642 assert_eq!(resolved_end, Some(end));
1643 }
1644
1645 #[rstest]
1646 fn test_resolve_request_time_range_makes_empty_interval_non_empty() {
1647 let end = UnixNanos::from(1_706_443_200_000_000_001);
1648
1649 let (start, resolved_end) = resolve_request_time_range(Some(end), Some(end));
1650
1651 assert_eq!(start, end - DurationNanos::new(1));
1652 assert_eq!(resolved_end, Some(end));
1653 }
1654
1655 #[rstest]
1656 fn test_requested_instrument_filters_exact_id() {
1657 let requested_id = InstrumentId::from("BTCUSDT.BINANCE");
1658 let instruments = vec![
1659 currency_pair("ETHUSDT.BINANCE"),
1660 currency_pair("BTCUSDT.BINANCE"),
1661 ];
1662
1663 let instrument = requested_instrument(instruments, requested_id).expect("instrument");
1664
1665 assert_eq!(instrument.id(), requested_id);
1666 }
1667
1668 #[rstest]
1669 fn test_requested_instrument_returns_latest_matching_id() {
1670 let requested_id = InstrumentId::from("BTCUSDT.BINANCE");
1671 let instruments = vec![
1672 currency_pair_with_ts_init("BTCUSDT.BINANCE", UnixNanos::from(1)),
1673 currency_pair_with_ts_init("BTCUSDT.BINANCE", UnixNanos::from(2)),
1674 ];
1675
1676 let instrument = requested_instrument(instruments, requested_id).expect("instrument");
1677
1678 assert_eq!(instrument.ts_init(), UnixNanos::from(2));
1679 }
1680
1681 #[rstest]
1682 fn test_requested_instrument_returns_none_on_miss() {
1683 let instruments = vec![currency_pair("ETHUSDT.BINANCE")];
1684
1685 let instrument = requested_instrument(instruments, InstrumentId::from("BTCUSDT.BINANCE"));
1686
1687 assert!(instrument.is_none());
1688 }
1689
1690 #[rstest]
1691 fn test_price_precision_from_params() {
1692 let mut params = Params::new();
1693 params.insert(PRICE_PRECISION_PARAM.to_string(), json!(5));
1694
1695 let price_precision = price_precision_from_params(Some(¶ms)).unwrap();
1696
1697 assert_eq!(price_precision, Some(5));
1698 }
1699
1700 #[rstest]
1701 fn test_price_precision_from_params_rejects_out_of_range_value() {
1702 let mut params = Params::new();
1703 params.insert(
1704 PRICE_PRECISION_PARAM.to_string(),
1705 json!(u64::from(u8::MAX) + 1),
1706 );
1707
1708 let result = price_precision_from_params(Some(¶ms));
1709
1710 assert!(result.is_err());
1711 }
1712
1713 #[rstest]
1714 fn test_partition_book_deltas_response_by_child_instrument() {
1715 let correlation_id = UUID4::new();
1716 let client_id = ClientId::from("DATABENTO-TEST");
1717 let parent = InstrumentId::from("ES.FUT.GLBX");
1718 let child_a = InstrumentId::from("ESM6.GLBX");
1719 let child_b = InstrumentId::from("ESU6.GLBX");
1720 let deltas = vec![
1721 OrderBookDelta::clear(child_a, 1, UnixNanos::from(1_000), UnixNanos::from(1_000)),
1722 OrderBookDelta::clear(child_b, 2, UnixNanos::from(2_000), UnixNanos::from(2_000)),
1723 OrderBookDelta::clear(child_a, 3, UnixNanos::from(3_000), UnixNanos::from(3_000)),
1724 ];
1725 let mut params = Params::new();
1726 params.insert(PRICE_PRECISION_PARAM.to_string(), json!(5));
1727 let response = BookDeltasResponse::new(
1728 correlation_id,
1729 client_id,
1730 parent,
1731 deltas.clone(),
1732 Some(UnixNanos::from(500)),
1733 Some(UnixNanos::from(4_000)),
1734 UnixNanos::from(5_000),
1735 Some(params.clone()),
1736 );
1737
1738 let responses = partition_book_deltas_response(response);
1739
1740 assert_eq!(responses.len(), 2);
1741 assert_eq!(responses[0].correlation_id, correlation_id);
1742 assert_eq!(responses[1].correlation_id, correlation_id);
1743 assert_eq!(responses[0].client_id, client_id);
1744 assert_eq!(responses[1].client_id, client_id);
1745 assert_eq!(responses[0].instrument_id, child_a);
1746 assert_eq!(responses[1].instrument_id, child_b);
1747 assert_eq!(responses[0].data, vec![deltas[0], deltas[2]]);
1748 assert_eq!(responses[1].data, vec![deltas[1]]);
1749 assert_eq!(responses[0].start, Some(UnixNanos::from(500)));
1750 assert_eq!(responses[1].start, Some(UnixNanos::from(500)));
1751 assert_eq!(responses[0].end, Some(UnixNanos::from(4_000)));
1752 assert_eq!(responses[1].end, Some(UnixNanos::from(4_000)));
1753 assert_eq!(responses[0].ts_init, UnixNanos::from(5_000));
1754 assert_eq!(responses[1].ts_init, UnixNanos::from(5_000));
1755 assert_eq!(responses[0].params, Some(params.clone()));
1756 assert_eq!(responses[1].params, Some(params));
1757 }
1758
1759 #[rstest]
1760 fn test_partition_book_deltas_response_preserves_homogeneous_response() {
1761 let correlation_id = UUID4::new();
1762 let client_id = ClientId::from("DATABENTO-TEST");
1763 let instrument_id = InstrumentId::from("ESM6.GLBX");
1764 let delta = OrderBookDelta::clear(
1765 instrument_id,
1766 1,
1767 UnixNanos::from(1_000),
1768 UnixNanos::from(1_000),
1769 );
1770 let response = BookDeltasResponse::new(
1771 correlation_id,
1772 client_id,
1773 instrument_id,
1774 vec![delta],
1775 Some(UnixNanos::from(500)),
1776 Some(UnixNanos::from(1_500)),
1777 UnixNanos::from(2_000),
1778 None,
1779 );
1780
1781 let responses = partition_book_deltas_response(response);
1782
1783 assert_eq!(responses.len(), 1);
1784 assert_eq!(responses[0].correlation_id, correlation_id);
1785 assert_eq!(responses[0].client_id, client_id);
1786 assert_eq!(responses[0].instrument_id, instrument_id);
1787 assert_eq!(responses[0].data, vec![delta]);
1788 assert_eq!(responses[0].start, Some(UnixNanos::from(500)));
1789 assert_eq!(responses[0].end, Some(UnixNanos::from(1_500)));
1790 assert_eq!(responses[0].ts_init, UnixNanos::from(2_000));
1791 assert_eq!(responses[0].params, None);
1792 }
1793
1794 #[rstest]
1795 fn test_partition_book_deltas_response_preserves_empty_parent_response() {
1796 let correlation_id = UUID4::new();
1797 let parent = InstrumentId::from("ES.FUT.GLBX");
1798 let response = BookDeltasResponse::new(
1799 correlation_id,
1800 ClientId::from("DATABENTO-TEST"),
1801 parent,
1802 Vec::new(),
1803 None,
1804 None,
1805 UnixNanos::from(1_000),
1806 None,
1807 );
1808
1809 let responses = partition_book_deltas_response(response);
1810
1811 assert_eq!(responses.len(), 1);
1812 let response = &responses[0];
1813 assert_eq!(response.correlation_id, correlation_id);
1814 assert_eq!(response.instrument_id, parent);
1815 assert!(response.data.is_empty());
1816 }
1817
1818 #[rstest]
1819 fn test_schema_from_params_returns_default() {
1820 let schema = schema_from_params(None, dbn::Schema::Mbp1, QUOTE_SCHEMAS).unwrap();
1821
1822 assert_eq!(schema, dbn::Schema::Mbp1);
1823 }
1824
1825 #[rstest]
1826 fn test_schema_from_params_accepts_allowed_value() {
1827 let mut params = Params::new();
1828 params.insert(SCHEMA_PARAM.to_string(), json!("tbbo"));
1829
1830 let schema = schema_from_params(Some(¶ms), dbn::Schema::Mbp1, QUOTE_SCHEMAS).unwrap();
1831
1832 assert_eq!(schema, dbn::Schema::Tbbo);
1833 }
1834
1835 #[rstest]
1836 fn test_schema_from_params_rejects_disallowed_value() {
1837 let mut params = Params::new();
1838 params.insert(SCHEMA_PARAM.to_string(), json!("mbo"));
1839
1840 let result = schema_from_params(Some(¶ms), dbn::Schema::Mbp1, QUOTE_SCHEMAS);
1841
1842 assert!(result.is_err());
1843 }
1844
1845 #[rstest]
1846 #[case::quotes(SubscribeKind::Quotes)]
1847 #[case::trades(SubscribeKind::Trades)]
1848 fn test_invalid_subscribe_params_do_not_create_feed_handler(#[case] kind: SubscribeKind) {
1849 let mut client = test_data_client();
1850 let mut params = Params::new();
1851 params.insert(SCHEMA_PARAM.to_string(), json!("definition"));
1852
1853 let result = match kind {
1854 SubscribeKind::Quotes => client.subscribe_quotes(subscribe_quotes_cmd(Some(params))),
1855 SubscribeKind::Trades => client.subscribe_trades(subscribe_trades_cmd(Some(params))),
1856 };
1857
1858 assert!(result.is_err());
1859 assert!(client.cmd_channels.lock().is_empty());
1860 }
1861
1862 #[rstest]
1863 fn test_send_subscription_commands_starts_after_subscribe() {
1864 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1865 let subscription = Subscription::builder()
1866 .schema(dbn::Schema::Mbp1)
1867 .symbols(vec!["ESM4"])
1868 .build();
1869
1870 send_subscription_commands(
1871 &tx,
1872 "GLBX.MDP3",
1873 Some((Symbol::from("ESM4"), 2)),
1874 subscription,
1875 true,
1876 )
1877 .unwrap();
1878
1879 assert!(matches!(
1880 rx.try_recv().unwrap(),
1881 HandlerCommand::SetPricePrecision(symbol, 2) if symbol == Symbol::from("ESM4")
1882 ));
1883 assert!(matches!(
1884 rx.try_recv().unwrap(),
1885 HandlerCommand::Subscribe(sub) if sub.schema == dbn::Schema::Mbp1
1886 ));
1887 assert!(matches!(rx.try_recv().unwrap(), HandlerCommand::Start));
1888 }
1889
1890 #[rstest]
1891 fn test_send_subscription_commands_without_precision_or_start() {
1892 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1893 let subscription = Subscription::builder()
1894 .schema(dbn::Schema::Mbp1)
1895 .symbols(vec!["ESM4"])
1896 .build();
1897
1898 send_subscription_commands(&tx, "GLBX.MDP3", None, subscription, false).unwrap();
1899
1900 assert!(matches!(
1901 rx.try_recv().unwrap(),
1902 HandlerCommand::Subscribe(sub) if sub.schema == dbn::Schema::Mbp1
1903 ));
1904 assert!(matches!(
1905 rx.try_recv(),
1906 Err(tokio::sync::mpsc::error::TryRecvError::Empty)
1907 ));
1908 }
1909}