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