Skip to main content

nautilus_databento/
data.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Provides a unified data client that combines Databento's live streaming and historical data capabilities.
17//!
18//! This module implements a data client that manages connections to multiple Databento datasets,
19//! handles live market data subscriptions, and provides access to historical data on demand.
20
21use 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/// Configuration for the Databento data client.
94#[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    /// Databento API credential.
105    pub(crate) credential: Credential,
106    /// Path to publishers.json file.
107    pub publishers_filepath: PathBuf,
108    /// Override for the Databento Historical API base URL.
109    pub historical_base_url: Option<String>,
110    /// Override for the Databento live gateway address.
111    pub live_gateway_addr: Option<String>,
112    /// Venue-to-dataset overrides applied on top of the publishers.json mappings.
113    pub venue_dataset_map: IndexMap<String, String>,
114    /// Whether to use exchange as venue for GLBX instruments.
115    pub use_exchange_as_venue: bool,
116    /// Whether to timestamp bars on close.
117    pub bars_timestamp_on_close: bool,
118    /// Reconnection timeout in minutes (None for infinite retries).
119    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    /// Creates a new [`DatabentoDataClientConfig`] instance.
134    #[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), // Default: 10 minutes
148            historical_base_url: None,
149            live_gateway_addr: None,
150        }
151    }
152
153    /// Returns the API key associated with this config.
154    #[must_use]
155    pub fn api_key(&self) -> &str {
156        self.credential.api_key()
157    }
158
159    /// Returns a masked version of the API key for logging purposes.
160    #[must_use]
161    pub fn api_key_masked(&self) -> String {
162        self.credential.api_key_masked()
163    }
164}
165
166/// A Databento data client that combines live streaming and historical data functionality.
167///
168/// This client uses the existing `DatabentoFeedHandler` for live data subscriptions
169/// and `DatabentoHistoricalClient` for historical data requests. It supports multiple
170/// datasets simultaneously, with separate feed handlers per dataset.
171#[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    /// Creates a new [`DatabentoDataClient`] instance.
193    ///
194    /// # Errors
195    ///
196    /// Returns an error if client creation or publisher configuration loading fails.
197    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        // Create data loader for venue-to-dataset mapping
219        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        // Load publisher configuration
228        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    /// Returns the API key associated with this client.
257    #[must_use]
258    pub fn api_key(&self) -> &str {
259        self.config.api_key()
260    }
261
262    /// Returns a masked version of the API key for logging purposes.
263    #[must_use]
264    pub fn api_key_masked(&self) -> String {
265        self.config.api_key_masked()
266    }
267
268    /// Gets the dataset for a given venue using the data loader.
269    ///
270    /// # Errors
271    ///
272    /// Returns an error if the venue-to-dataset mapping cannot be found.
273    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    /// Gets or creates a feed handler for the specified dataset.
281    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    /// Initializes the live feed handler for streaming data.
348    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        // Spawn message processing task with cancellation support
384        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    /// Returns the client identifier.
464    fn client_id(&self) -> ClientId {
465        self.client_id
466    }
467
468    /// Returns the venue associated with this client (None for multi-venue clients).
469    fn venue(&self) -> Option<Venue> {
470        None
471    }
472
473    /// Starts the data client.
474    ///
475    /// # Errors
476    ///
477    /// Returns an error if the client fails to start.
478    fn start(&mut self) -> anyhow::Result<()> {
479        log::debug!("Starting");
480        Ok(())
481    }
482
483    /// Stops the data client and cancels all active subscriptions.
484    ///
485    /// # Errors
486    ///
487    /// Returns an error if the client fails to stop cleanly.
488    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    /// Returns whether the client is currently connected.
554    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    /// Subscribes to instrument definition data for the specified instrument.
563    ///
564    /// # Errors
565    ///
566    /// Returns an error if the subscription request fails.
567    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    /// Subscribes to quote tick data for the specified instruments.
586    ///
587    /// # Errors
588    ///
589    /// Returns an error if the subscription request fails.
590    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    /// Subscribes to trade tick data for the specified instruments.
617    ///
618    /// # Errors
619    ///
620    /// Returns an error if the subscription request fails.
621    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    /// Subscribes to order book delta updates for the specified instruments.
648    ///
649    /// # Errors
650    ///
651    /// Returns an error if the subscription request fails.
652    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) // Market by order for book deltas
662            .symbols(symbol)
663            .build();
664
665        self.send_subscription_to_dataset(&dataset, None, subscription, start_after_subscribe)?;
666
667        Ok(())
668    }
669
670    /// Subscribes to instrument status updates for the specified instruments.
671    ///
672    /// # Errors
673    ///
674    /// Returns an error if the subscription request fails.
675    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    // Unsubscribe methods
697    fn unsubscribe_quotes(&mut self, cmd: &UnsubscribeQuotes) -> anyhow::Result<()> {
698        // Note: Databento live API doesn't support granular unsubscribing.
699        // The feed handler manages subscriptions and can handle reconnections
700        // with the appropriate subscription state.
701        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        // Note: Databento live API doesn't support granular unsubscribing.
711        // The feed handler manages subscriptions and can handle reconnections
712        // with the appropriate subscription state.
713        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        // Note: Databento live API doesn't support granular unsubscribing.
723        // The feed handler manages subscriptions and can handle reconnections
724        // with the appropriate subscription state.
725        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        // Note: Databento live API doesn't support granular unsubscribing.
738        // The feed handler manages subscriptions and can handle reconnections
739        // with the appropriate subscription state.
740        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")] // overrides the apply_default EQUS -> EQUS.MINI mapping
1544    #[case("GLBX", "EQUS.MINI")] // overrides the apply_default GLBX -> GLBX.MDP3 mapping
1545    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        // The override is targeted: an unrelated venue keeps its default.
1570        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(&params)).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(&params));
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(&params), 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(&params), 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}