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