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