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