Skip to main content

nautilus_data/engine/
streaming.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
16use ahash::AHashMap;
17use jiff::Timestamp;
18use nautilus_common::messages::data::{
19    BarsResponse, BookDeltasResponse, BookDepthResponse, CustomDataResponse, DataResponse,
20    FundingRatesResponse, InstrumentResponse, InstrumentsResponse, QuotesResponse, RequestBars,
21    RequestBookDeltas, RequestBookDepth, RequestCommand, RequestCustomData, RequestFundingRates,
22    RequestInstrument, RequestInstruments, RequestQuotes, RequestTrades, SubscribeBars,
23    SubscribeCommand, SubscribeCustomData, SubscribeQuotes, SubscribeTrades, TradesResponse,
24};
25use nautilus_core::{
26    Params, UUID4, UnixNanos,
27    correctness::{FAILED, check_key_not_in_map},
28};
29use nautilus_model::{
30    data::{
31        Bar, CustomData, Data, FundingRateUpdate, OrderBookDelta, OrderBookDepth10, QuoteTick,
32        TradeTick,
33    },
34    identifiers::{ClientId, Venue},
35    instruments::{Instrument, InstrumentAny},
36};
37use nautilus_persistence::backend::catalog::ParquetDataCatalog;
38use serde_json::Value;
39use ustr::Ustr;
40
41use super::{DataEngine, requests::request_params};
42
43const PARAM_SKIP_CATALOG_DATA: &str = "skip_catalog_data";
44const PARAM_UPDATE_CATALOG: &str = "update_catalog";
45const PARAM_FORCE_INSTRUMENT_UPDATE: &str = "force_instrument_update";
46const PARAM_SUBSCRIPTION_NAME: &str = "subscription_name";
47const PARAM_FROM_DAY_START: &str = "from_day_start";
48const CATALOG_CLIENT_ID: &str = "CATALOG";
49
50pub(crate) type CatalogMap = AHashMap<Ustr, ParquetDataCatalog>;
51
52impl DataEngine {
53    /// Registers the `catalog` with the engine with an optional specific `name`.
54    ///
55    /// # Panics
56    ///
57    /// Panics if a catalog with the same `name` has already been registered.
58    pub fn register_catalog(&mut self, catalog: ParquetDataCatalog, name: Option<&str>) {
59        let name = Ustr::from(name.unwrap_or("catalog_0"));
60
61        check_key_not_in_map(&name, &self.catalogs, "name", "catalogs").expect(FAILED);
62
63        self.catalogs.insert(name, catalog);
64        log::info!("Registered catalog <{name}>");
65    }
66
67    pub(super) fn subscribe_command_with_prefilled_start_ns(
68        &self,
69        cmd: SubscribeCommand,
70    ) -> anyhow::Result<SubscribeCommand> {
71        match cmd {
72            SubscribeCommand::Quotes(cmd) if Self::is_start_ns_missing(cmd.params.as_ref()) => {
73                let identifier = cmd.instrument_id.to_string();
74                let params = self.params_with_prefilled_start_ns(
75                    cmd.params.as_ref(),
76                    "quotes",
77                    &identifier,
78                )?;
79                Ok(SubscribeCommand::Quotes(SubscribeQuotes { params, ..cmd }))
80            }
81            SubscribeCommand::Trades(cmd) if Self::is_start_ns_missing(cmd.params.as_ref()) => {
82                let identifier = cmd.instrument_id.to_string();
83                let params = self.params_with_prefilled_start_ns(
84                    cmd.params.as_ref(),
85                    "trades",
86                    &identifier,
87                )?;
88                Ok(SubscribeCommand::Trades(SubscribeTrades { params, ..cmd }))
89            }
90            SubscribeCommand::Bars(cmd)
91                if cmd.bar_type.is_externally_aggregated()
92                    && Self::is_start_ns_missing(cmd.params.as_ref()) =>
93            {
94                let identifier = cmd.bar_type.to_string();
95                let params =
96                    self.params_with_prefilled_start_ns(cmd.params.as_ref(), "bars", &identifier)?;
97                Ok(SubscribeCommand::Bars(SubscribeBars { params, ..cmd }))
98            }
99            SubscribeCommand::Data(cmd) if Self::is_start_ns_missing(cmd.params.as_ref()) => {
100                let type_name = cmd.data_type.type_name().to_string();
101                let identifier = cmd.data_type.identifier().map(String::from);
102                let params = self.params_with_custom_data_prefilled_start_ns(
103                    cmd.params.as_ref(),
104                    &type_name,
105                    identifier.as_deref(),
106                )?;
107                Ok(SubscribeCommand::Data(SubscribeCustomData {
108                    params,
109                    ..cmd
110                }))
111            }
112            _ => Ok(cmd),
113        }
114    }
115
116    fn is_start_ns_missing(params: Option<&Params>) -> bool {
117        params.is_none_or(|params| !params.contains_key("start_ns"))
118    }
119
120    fn params_with_prefilled_start_ns(
121        &self,
122        params: Option<&Params>,
123        data_cls: &str,
124        identifier: &str,
125    ) -> anyhow::Result<Option<Params>> {
126        let last_timestamp = self.catalog_last_timestamp(data_cls, identifier)?;
127
128        Ok(Some(Self::params_with_start_ns(params, last_timestamp)))
129    }
130
131    fn params_with_custom_data_prefilled_start_ns(
132        &self,
133        params: Option<&Params>,
134        type_name: &str,
135        identifier: Option<&str>,
136    ) -> anyhow::Result<Option<Params>> {
137        let last_timestamp = self.catalog_custom_data_last_timestamp(type_name, identifier)?;
138
139        Ok(Some(Self::params_with_start_ns(params, last_timestamp)))
140    }
141
142    fn params_with_start_ns(params: Option<&Params>, last_timestamp: Option<u64>) -> Params {
143        let start_ns = last_timestamp.map_or(Value::Null, |last_timestamp| {
144            Value::from(last_timestamp.saturating_add(1))
145        });
146        let mut params = params.cloned().unwrap_or_else(Params::new);
147
148        params.insert("start_ns".to_string(), start_ns);
149
150        params
151    }
152
153    fn catalog_last_timestamp(
154        &self,
155        data_cls: &str,
156        identifier: &str,
157    ) -> anyhow::Result<Option<u64>> {
158        for catalog in self.catalogs.values() {
159            if let Some(last_timestamp) =
160                catalog.query_last_timestamp(data_cls, Some(identifier))?
161            {
162                return Ok(Some(last_timestamp));
163            }
164        }
165
166        Ok(None)
167    }
168
169    fn catalog_custom_data_last_timestamp(
170        &self,
171        type_name: &str,
172        identifier: Option<&str>,
173    ) -> anyhow::Result<Option<u64>> {
174        for catalog in self.catalogs.values() {
175            let last_timestamp = if let Some(identifier) = identifier {
176                let directory = catalog.make_path_custom_data(type_name, Some(identifier))?;
177                let intervals = catalog.get_directory_intervals(&directory)?;
178                intervals.last().map(|(_, last_timestamp)| *last_timestamp)
179            } else {
180                let data_cls = format!("custom/{type_name}");
181                catalog.query_last_timestamp(&data_cls, None)?
182            };
183
184            if let Some(last_timestamp) = last_timestamp {
185                return Ok(Some(last_timestamp));
186            }
187        }
188
189        Ok(None)
190    }
191
192    pub(super) fn catalogs_registered(&self) -> bool {
193        !self.catalogs.is_empty()
194    }
195
196    // Bounds the request window, walks the catalogs to find one whose missing
197    // intervals differ from the full requested range, then fans the parent out via
198    // the pipeline with one catalog leg plus one client leg per missing interval.
199    // With no catalog match and no resolvable client, the engine emits an empty
200    // response keyed by the parent request ID.
201    pub(super) fn dispatch_date_range_request(
202        &mut self,
203        req: RequestCommand,
204    ) -> anyhow::Result<()> {
205        if matches!(
206            req,
207            RequestCommand::Instrument(_) | RequestCommand::Instruments(_)
208        ) {
209            return self.dispatch_instrument_catalog_request(req);
210        }
211
212        let Some(key) = request_identifier(&req) else {
213            return self.dispatch_request_to_client(req).map(|_| ());
214        };
215
216        let now_ns = self.clock.borrow().timestamp_ns();
217        let now_dt = now_ns.to_datetime_utc();
218        let query_past_data = request_params(&req)
219            .and_then(|p| p.get(PARAM_SUBSCRIPTION_NAME))
220            .is_none();
221
222        let (start_dt, end_dt) = bound_request_dates(
223            request_start(&req),
224            request_end(&req),
225            now_dt,
226            query_past_data,
227        );
228        let start_ns = datetime_to_unix_nanos_or_zero(start_dt);
229        let end_ns = datetime_to_unix_nanos_or_zero(end_dt);
230
231        if start_ns > end_ns {
232            anyhow::bail!(
233                "Cannot dispatch request, start {start_ns} was greater than end {end_ns}"
234            );
235        }
236
237        let client_id = req.client_id().copied();
238        let venue = req.venue().copied();
239        let used_client_id = self
240            .get_client(client_id.as_ref(), venue.as_ref())
241            .map(|client| client.client_id());
242
243        // Floor the catalog window to the UTC day boundary so the day-start F_SNAPSHOT frame is
244        // selected and read for the snapshot replay; client gaps keep the original window.
245        // The parent request keeps its original start, so the merged response trims back to it.
246        let (catalog_start_dt, catalog_start_ns) = if matches!(req, RequestCommand::BookDeltas(_))
247            && request_params(&req)
248                .and_then(|p| p.get_bool(PARAM_FROM_DAY_START))
249                .unwrap_or(true)
250        {
251            let floored = floor_to_utc_day(start_dt);
252            (floored, datetime_to_unix_nanos_or_zero(floored))
253        } else {
254            (start_dt, start_ns)
255        };
256
257        let query_interval = vec![(start_ns.as_u64(), end_ns.as_u64())];
258        let catalog_query_interval = vec![(catalog_start_ns.as_u64(), end_ns.as_u64())];
259        let mut missing_intervals = query_interval.clone();
260        let mut has_catalog_data = false;
261        let mut winning_catalog: Option<Ustr> = None;
262
263        for (name, catalog) in &self.catalogs {
264            let catalog_intervals = catalog_missing_intervals(
265                catalog,
266                catalog_start_ns.as_u64(),
267                end_ns.as_u64(),
268                &key,
269            )?;
270
271            if catalog_intervals != catalog_query_interval {
272                has_catalog_data = true;
273                winning_catalog = Some(*name);
274                // Client legs fill only the requested window, not the pre-start range
275                missing_intervals = if catalog_start_ns == start_ns {
276                    catalog_intervals
277                } else {
278                    catalog_missing_intervals(catalog, start_ns.as_u64(), end_ns.as_u64(), &key)?
279                };
280                break;
281            }
282        }
283
284        let skip_catalog_data = request_params(&req)
285            .and_then(|p| p.get_bool(PARAM_SKIP_CATALOG_DATA))
286            .unwrap_or(false);
287
288        // When `skip_catalog_data` is set the client must serve the full parent window;
289        // dropping the catalog leg without resetting the missing intervals would leave
290        // the catalog-covered range unanswered.
291        if skip_catalog_data {
292            missing_intervals = query_interval;
293        }
294
295        let n_client_requests = if used_client_id.is_some() {
296            missing_intervals.len()
297        } else {
298            0
299        };
300        let n_catalog_requests = usize::from(has_catalog_data && !skip_catalog_data);
301        let n_requests = n_client_requests + n_catalog_requests;
302
303        if n_requests == 0 {
304            let empty = build_empty_response(&req, start_ns, end_ns, used_client_id, now_ns)?;
305            self.response(empty);
306            return Ok(());
307        }
308
309        let parent_id = *req.request_id();
310        self.new_request_pipeline(req.clone(), n_requests);
311
312        if n_catalog_requests == 1
313            && let Some(catalog_name) = winning_catalog
314        {
315            let leg = with_dates_for_pipeline(&req, Some(catalog_start_dt), Some(end_dt), now_ns);
316            let leg_id = *leg.request_id();
317            self.register_request_pipeline_leg(leg_id, parent_id);
318
319            match self.query_catalog_leg(
320                &leg,
321                catalog_name,
322                catalog_start_ns,
323                end_ns,
324                used_client_id,
325                now_ns,
326            ) {
327                Ok(resp) => self.response(resp),
328                Err(e) => {
329                    log::error!(
330                        "Catalog leg query failed for parent {parent_id} (catalog {catalog_name}): {e}"
331                    );
332                    let empty = match build_empty_response(
333                        &leg,
334                        start_ns,
335                        end_ns,
336                        used_client_id,
337                        now_ns,
338                    ) {
339                        Ok(empty) => empty,
340                        Err(e) => {
341                            self.abort_request_pipeline(parent_id);
342                            return Err(e);
343                        }
344                    };
345                    self.response(empty);
346                }
347            }
348        }
349
350        if n_client_requests > 0 {
351            for (leg_start_ns, leg_end_ns) in &missing_intervals {
352                let leg_start_dt = UnixNanos::from(*leg_start_ns).to_datetime_utc();
353                let leg_end_dt = UnixNanos::from(*leg_end_ns).to_datetime_utc();
354                let leg =
355                    with_dates_for_pipeline(&req, Some(leg_start_dt), Some(leg_end_dt), now_ns);
356                let leg_id = *leg.request_id();
357                self.register_request_pipeline_leg(leg_id, parent_id);
358
359                if let Err(e) = self.dispatch_request_to_client(leg) {
360                    // Abort the whole pipeline so the parent does not stay half-registered
361                    // waiting on a leg the client never accepted. Any catalog leg already
362                    // buffered for this parent is discarded with the pipeline state.
363                    log::error!("Client leg dispatch failed for parent {parent_id}: {e}");
364                    self.abort_request_pipeline(parent_id);
365                    return Err(e);
366                }
367            }
368        }
369
370        Ok(())
371    }
372
373    fn abort_request_pipeline(&mut self, parent_id: UUID4) {
374        self.request_pipeline_n_components.remove(&parent_id);
375        self.request_pipeline_parent_request.remove(&parent_id);
376        self.request_pipeline_responses.remove(&parent_id);
377        self.request_pipeline_parent_request_id
378            .retain(|_, p_id| *p_id != parent_id);
379    }
380
381    fn query_catalog_leg(
382        &mut self,
383        leg: &RequestCommand,
384        catalog_name: Ustr,
385        start_ns: UnixNanos,
386        end_ns: UnixNanos,
387        used_client_id: Option<ClientId>,
388        ts_init: UnixNanos,
389    ) -> anyhow::Result<DataResponse> {
390        let catalog = self.catalogs.get_mut(&catalog_name).ok_or_else(|| {
391            anyhow::anyhow!("Catalog {catalog_name} disappeared between intervals query and read")
392        })?;
393
394        match leg {
395            RequestCommand::Quotes(cmd) => {
396                let data: Vec<QuoteTick> = catalog.quote_ticks(
397                    Some(vec![cmd.instrument_id.to_string()]),
398                    Some(start_ns),
399                    Some(end_ns),
400                )?;
401                Ok(build_quotes_catalog_response(
402                    cmd,
403                    data,
404                    start_ns,
405                    end_ns,
406                    used_client_id,
407                    ts_init,
408                ))
409            }
410            RequestCommand::Trades(cmd) => {
411                let data: Vec<TradeTick> = catalog.trade_ticks(
412                    Some(vec![cmd.instrument_id.to_string()]),
413                    Some(start_ns),
414                    Some(end_ns),
415                )?;
416                Ok(build_trades_catalog_response(
417                    cmd,
418                    data,
419                    start_ns,
420                    end_ns,
421                    used_client_id,
422                    ts_init,
423                ))
424            }
425            RequestCommand::FundingRates(cmd) => {
426                let data: Vec<FundingRateUpdate> = catalog.funding_rates(
427                    Some(vec![cmd.instrument_id.to_string()]),
428                    Some(start_ns),
429                    Some(end_ns),
430                )?;
431                Ok(build_funding_rates_catalog_response(
432                    cmd,
433                    data,
434                    start_ns,
435                    end_ns,
436                    used_client_id,
437                    ts_init,
438                ))
439            }
440            RequestCommand::Bars(cmd) => {
441                let data: Vec<Bar> = catalog.bars(
442                    Some(vec![cmd.bar_type.to_string()]),
443                    Some(start_ns),
444                    Some(end_ns),
445                )?;
446                Ok(build_bars_catalog_response(
447                    cmd,
448                    data,
449                    start_ns,
450                    end_ns,
451                    used_client_id,
452                    ts_init,
453                ))
454            }
455            RequestCommand::Data(cmd) => {
456                let identifiers = cmd
457                    .data_type
458                    .identifier()
459                    .map(|identifier| vec![identifier.to_string()]);
460                let where_clause = cmd
461                    .params
462                    .as_ref()
463                    .and_then(|params| params.get_str("filter_expr"));
464                let data = catalog.query_custom_data_dynamic(
465                    cmd.data_type.type_name(),
466                    identifiers.as_deref(),
467                    Some(start_ns),
468                    Some(end_ns),
469                    where_clause,
470                    None,
471                    true,
472                )?;
473                Ok(build_custom_data_catalog_response(
474                    cmd,
475                    custom_data_from_dynamic(data),
476                    start_ns,
477                    end_ns,
478                    ts_init,
479                ))
480            }
481            RequestCommand::BookDeltas(cmd) => {
482                let data: Vec<OrderBookDelta> = catalog.order_book_deltas(
483                    Some(vec![cmd.instrument_id.to_string()]),
484                    Some(start_ns),
485                    Some(end_ns),
486                )?;
487                Ok(build_book_deltas_catalog_response(
488                    cmd,
489                    data,
490                    start_ns,
491                    end_ns,
492                    used_client_id,
493                    ts_init,
494                ))
495            }
496            RequestCommand::BookDepth(cmd) => {
497                let data: Vec<OrderBookDepth10> = catalog.order_book_depth10(
498                    Some(vec![cmd.instrument_id.to_string()]),
499                    Some(start_ns),
500                    Some(end_ns),
501                )?;
502                Ok(build_book_depth_catalog_response(
503                    cmd,
504                    data,
505                    start_ns,
506                    end_ns,
507                    used_client_id,
508                    ts_init,
509                ))
510            }
511            _ => {
512                anyhow::bail!("query_catalog_leg called with non-catalog-eligible variant {leg:?}")
513            }
514        }
515    }
516
517    fn dispatch_instrument_catalog_request(&mut self, req: RequestCommand) -> anyhow::Result<()> {
518        match req {
519            RequestCommand::Instrument(cmd) => self.dispatch_instrument_request(cmd),
520            RequestCommand::Instruments(cmd) => self.dispatch_instruments_request(cmd),
521            _ => self.dispatch_request_to_client(req).map(|_| ()),
522        }
523    }
524
525    fn dispatch_instrument_request(&mut self, cmd: RequestInstrument) -> anyhow::Result<()> {
526        let force_instrument_update = cmd
527            .params
528            .as_ref()
529            .and_then(|params| params.get_bool(PARAM_FORCE_INSTRUMENT_UPDATE))
530            .unwrap_or(false);
531
532        if force_instrument_update {
533            return self
534                .dispatch_request_to_client(RequestCommand::Instrument(cmd))
535                .map(|_| ());
536        }
537
538        let identifier = cmd.instrument_id.to_string();
539        let Some(catalog_name) = self.catalog_with_last_timestamp("instruments", &identifier)?
540        else {
541            return self
542                .dispatch_request_to_client(RequestCommand::Instrument(cmd))
543                .map(|_| ());
544        };
545
546        let now_ns = self.clock.borrow().timestamp_ns();
547        let used_client_id = self
548            .get_client(cmd.client_id.as_ref(), Some(&cmd.instrument_id.venue))
549            .map(|client| client.client_id());
550        let (start_dt, end_dt) =
551            bound_request_dates(cmd.start, cmd.end, now_ns.to_datetime_utc(), true);
552        let start_ns = datetime_to_unix_nanos_or_zero(start_dt);
553        let end_ns = datetime_to_unix_nanos_or_zero(end_dt);
554        let query_end = cmd.end.map(datetime_to_unix_nanos_or_zero);
555        let catalog = self.catalogs.get(&catalog_name).ok_or_else(|| {
556            anyhow::anyhow!("Catalog {catalog_name} disappeared between timestamp query and read")
557        })?;
558        let mut data = catalog.instruments(
559            Some(std::slice::from_ref(&identifier)),
560            Some(start_ns),
561            query_end,
562        )?;
563        data = latest_instruments(data);
564
565        if let Some(instrument) = data.into_iter().next() {
566            let response = DataResponse::Instrument(Box::new(InstrumentResponse::new(
567                cmd.request_id,
568                resolve_response_client_id(cmd.client_id, used_client_id),
569                cmd.instrument_id,
570                instrument,
571                Some(start_ns),
572                Some(end_ns),
573                now_ns,
574                Some(catalog_response_params(cmd.params.as_ref())),
575            )));
576            self.response(response);
577            return Ok(());
578        }
579
580        self.dispatch_request_to_client(RequestCommand::Instrument(cmd))
581            .map(|_| ())
582    }
583
584    fn dispatch_instruments_request(&mut self, cmd: RequestInstruments) -> anyhow::Result<()> {
585        let update_catalog = cmd
586            .params
587            .as_ref()
588            .and_then(|params| params.get_bool(PARAM_UPDATE_CATALOG))
589            .unwrap_or(false);
590        let force_instrument_update = cmd
591            .params
592            .as_ref()
593            .and_then(|params| params.get_bool(PARAM_FORCE_INSTRUMENT_UPDATE))
594            .unwrap_or(false);
595
596        if update_catalog || force_instrument_update {
597            return self
598                .dispatch_request_to_client(RequestCommand::Instruments(cmd))
599                .map(|_| ());
600        }
601
602        let now_ns = self.clock.borrow().timestamp_ns();
603        let used_client_id = self
604            .get_client(cmd.client_id.as_ref(), cmd.venue.as_ref())
605            .map(|client| client.client_id());
606        let (start_dt, end_dt) =
607            bound_request_dates(cmd.start, cmd.end, now_ns.to_datetime_utc(), true);
608        let start_ns = datetime_to_unix_nanos_or_zero(start_dt);
609        let end_ns = datetime_to_unix_nanos_or_zero(end_dt);
610        let query_end = cmd.end.map(datetime_to_unix_nanos_or_zero);
611        let mut data = Vec::new();
612
613        for catalog in self.catalogs.values() {
614            data.extend(catalog.instruments(None, Some(start_ns), query_end)?);
615        }
616
617        if let Some(venue) = cmd.venue {
618            data.retain(|instrument| instrument.venue() == venue);
619        }
620
621        if instrument_only_last(cmd.params.as_ref()) {
622            data = latest_instruments(data);
623        }
624
625        let response = DataResponse::Instruments(InstrumentsResponse::new(
626            cmd.request_id,
627            resolve_response_client_id(cmd.client_id, used_client_id),
628            instrument_response_venue(cmd.venue, &data),
629            data,
630            Some(start_ns),
631            Some(end_ns),
632            now_ns,
633            Some(catalog_response_params(cmd.params.as_ref())),
634        ));
635        self.response(response);
636        Ok(())
637    }
638
639    fn catalog_with_last_timestamp(
640        &self,
641        data_cls: &str,
642        identifier: &str,
643    ) -> anyhow::Result<Option<Ustr>> {
644        for (name, catalog) in &self.catalogs {
645            if catalog
646                .query_last_timestamp(data_cls, Some(identifier))?
647                .is_some()
648            {
649                return Ok(Some(*name));
650            }
651        }
652
653        Ok(None)
654    }
655}
656
657struct RequestCatalogKey {
658    data_cls: String,
659    type_name: Option<String>,
660    identifier: Option<String>,
661}
662
663pub(super) fn is_date_range_variant(req: &RequestCommand) -> bool {
664    matches!(
665        req,
666        RequestCommand::Data(_)
667            | RequestCommand::Instrument(_)
668            | RequestCommand::Instruments(_)
669            | RequestCommand::Quotes(_)
670            | RequestCommand::Trades(_)
671            | RequestCommand::FundingRates(_)
672            | RequestCommand::Bars(_)
673            | RequestCommand::BookDeltas(_)
674            | RequestCommand::BookDepth(_)
675    )
676}
677
678fn request_identifier(req: &RequestCommand) -> Option<RequestCatalogKey> {
679    match req {
680        RequestCommand::Data(cmd) => Some(RequestCatalogKey {
681            data_cls: format!("custom/{}", cmd.data_type.type_name()),
682            type_name: Some(cmd.data_type.type_name().to_string()),
683            identifier: cmd.data_type.identifier().map(String::from),
684        }),
685        RequestCommand::Quotes(cmd) => Some(RequestCatalogKey::new(
686            "quotes",
687            Some(cmd.instrument_id.to_string()),
688        )),
689        RequestCommand::Trades(cmd) => Some(RequestCatalogKey::new(
690            "trades",
691            Some(cmd.instrument_id.to_string()),
692        )),
693        RequestCommand::FundingRates(cmd) => Some(RequestCatalogKey::new(
694            "funding_rate_update",
695            Some(cmd.instrument_id.to_string()),
696        )),
697        RequestCommand::Bars(cmd) => Some(RequestCatalogKey::new(
698            "bars",
699            Some(cmd.bar_type.to_string()),
700        )),
701        RequestCommand::BookDeltas(cmd) => Some(RequestCatalogKey::new(
702            "order_book_deltas",
703            Some(cmd.instrument_id.to_string()),
704        )),
705        RequestCommand::BookDepth(cmd) => Some(RequestCatalogKey::new(
706            "order_book_depths",
707            Some(cmd.instrument_id.to_string()),
708        )),
709        _ => None,
710    }
711}
712
713impl RequestCatalogKey {
714    fn new(data_cls: &str, identifier: Option<String>) -> Self {
715        Self {
716            data_cls: data_cls.to_string(),
717            type_name: None,
718            identifier,
719        }
720    }
721}
722
723fn catalog_missing_intervals(
724    catalog: &ParquetDataCatalog,
725    start: u64,
726    end: u64,
727    key: &RequestCatalogKey,
728) -> anyhow::Result<Vec<(u64, u64)>> {
729    if let Some(type_name) = key.type_name.as_deref()
730        && let Some(identifier) = key.identifier.as_deref()
731    {
732        let directory = catalog.make_path_custom_data(type_name, Some(identifier))?;
733        let intervals = catalog.get_directory_intervals(&directory)?;
734        return Ok(missing_interval_diff(start, end, &intervals));
735    }
736
737    catalog.get_missing_intervals_for_request(start, end, &key.data_cls, key.identifier.as_deref())
738}
739
740fn request_start(req: &RequestCommand) -> Option<Timestamp> {
741    match req {
742        RequestCommand::Data(cmd) => cmd.start,
743        RequestCommand::Instrument(cmd) => cmd.start,
744        RequestCommand::Instruments(cmd) => cmd.start,
745        RequestCommand::Quotes(cmd) => cmd.start,
746        RequestCommand::Trades(cmd) => cmd.start,
747        RequestCommand::FundingRates(cmd) => cmd.start,
748        RequestCommand::Bars(cmd) => cmd.start,
749        RequestCommand::BookDeltas(cmd) => cmd.start,
750        RequestCommand::BookDepth(cmd) => cmd.start,
751        _ => None,
752    }
753}
754
755fn request_end(req: &RequestCommand) -> Option<Timestamp> {
756    match req {
757        RequestCommand::Data(cmd) => cmd.end,
758        RequestCommand::Instrument(cmd) => cmd.end,
759        RequestCommand::Instruments(cmd) => cmd.end,
760        RequestCommand::Quotes(cmd) => cmd.end,
761        RequestCommand::Trades(cmd) => cmd.end,
762        RequestCommand::FundingRates(cmd) => cmd.end,
763        RequestCommand::Bars(cmd) => cmd.end,
764        RequestCommand::BookDeltas(cmd) => cmd.end,
765        RequestCommand::BookDepth(cmd) => cmd.end,
766        _ => None,
767    }
768}
769
770fn bound_request_dates(
771    start: Option<Timestamp>,
772    end: Option<Timestamp>,
773    now: Timestamp,
774    query_past_data: bool,
775) -> (Timestamp, Timestamp) {
776    let zero = Timestamp::UNIX_EPOCH;
777    let mut start = start.unwrap_or(zero);
778    let mut end = end.unwrap_or(now);
779
780    if query_past_data {
781        if start > now {
782            start = now;
783        }
784
785        if end > now {
786            end = now;
787        }
788    }
789
790    (start, end)
791}
792
793fn datetime_to_unix_nanos_or_zero(dt: Timestamp) -> UnixNanos {
794    UnixNanos::from(u64::try_from(dt.as_nanosecond().max(0)).unwrap_or(0))
795}
796
797fn floor_to_utc_day(dt: Timestamp) -> Timestamp {
798    let midnight = jiff::tz::Offset::UTC.to_datetime(dt).date().at(0, 0, 0, 0);
799    jiff::tz::Offset::UTC
800        .to_timestamp(midnight)
801        .expect("midnight UTC is always valid")
802}
803
804fn with_dates_for_pipeline(
805    req: &RequestCommand,
806    start: Option<Timestamp>,
807    end: Option<Timestamp>,
808    ts_init: UnixNanos,
809) -> RequestCommand {
810    let new_id = UUID4::new();
811
812    match req {
813        RequestCommand::Quotes(cmd) => RequestCommand::Quotes(RequestQuotes {
814            instrument_id: cmd.instrument_id,
815            start,
816            end,
817            limit: cmd.limit,
818            client_id: cmd.client_id,
819            request_id: new_id,
820            ts_init,
821            params: cmd.params.clone(),
822        }),
823        RequestCommand::Trades(cmd) => RequestCommand::Trades(RequestTrades {
824            instrument_id: cmd.instrument_id,
825            start,
826            end,
827            limit: cmd.limit,
828            client_id: cmd.client_id,
829            request_id: new_id,
830            ts_init,
831            params: cmd.params.clone(),
832        }),
833        RequestCommand::FundingRates(cmd) => RequestCommand::FundingRates(RequestFundingRates {
834            instrument_id: cmd.instrument_id,
835            start,
836            end,
837            limit: cmd.limit,
838            client_id: cmd.client_id,
839            request_id: new_id,
840            ts_init,
841            params: cmd.params.clone(),
842        }),
843        RequestCommand::BookDeltas(cmd) => RequestCommand::BookDeltas(RequestBookDeltas {
844            instrument_id: cmd.instrument_id,
845            start,
846            end,
847            limit: cmd.limit,
848            client_id: cmd.client_id,
849            request_id: new_id,
850            ts_init,
851            params: cmd.params.clone(),
852        }),
853        RequestCommand::BookDepth(cmd) => RequestCommand::BookDepth(RequestBookDepth {
854            instrument_id: cmd.instrument_id,
855            start,
856            end,
857            limit: cmd.limit,
858            depth: cmd.depth,
859            client_id: cmd.client_id,
860            request_id: new_id,
861            ts_init,
862            params: cmd.params.clone(),
863        }),
864        RequestCommand::Data(cmd) => RequestCommand::Data(RequestCustomData {
865            client_id: cmd.client_id,
866            data_type: cmd.data_type.clone(),
867            start,
868            end,
869            limit: cmd.limit,
870            request_id: new_id,
871            ts_init,
872            params: cmd.params.clone(),
873        }),
874        RequestCommand::Bars(cmd) => RequestCommand::Bars(RequestBars {
875            bar_type: cmd.bar_type,
876            start,
877            end,
878            limit: cmd.limit,
879            client_id: cmd.client_id,
880            request_id: new_id,
881            ts_init,
882            params: cmd.params.clone(),
883        }),
884        // `Join` and the non-date-range variants should never reach this path; the dispatcher
885        // gates on `is_date_range_variant` first. Cloning preserves behaviour if a caller
886        // reaches this arm.
887        _ => req.clone(),
888    }
889}
890
891fn build_empty_response(
892    req: &RequestCommand,
893    start: UnixNanos,
894    end: UnixNanos,
895    used_client_id: Option<ClientId>,
896    ts_init: UnixNanos,
897) -> anyhow::Result<DataResponse> {
898    let response = match req {
899        RequestCommand::Data(cmd) => DataResponse::Data(CustomDataResponse::new(
900            cmd.request_id,
901            cmd.client_id,
902            None,
903            cmd.data_type.clone(),
904            Vec::<CustomData>::new(),
905            Some(start),
906            Some(end),
907            ts_init,
908            cmd.params.clone(),
909        )),
910        RequestCommand::Quotes(cmd) => DataResponse::Quotes(QuotesResponse::new(
911            cmd.request_id,
912            resolve_response_client_id(cmd.client_id, used_client_id),
913            cmd.instrument_id,
914            Vec::new(),
915            Some(start),
916            Some(end),
917            ts_init,
918            cmd.params.clone(),
919        )),
920        RequestCommand::Trades(cmd) => DataResponse::Trades(TradesResponse::new(
921            cmd.request_id,
922            resolve_response_client_id(cmd.client_id, used_client_id),
923            cmd.instrument_id,
924            Vec::new(),
925            Some(start),
926            Some(end),
927            ts_init,
928            cmd.params.clone(),
929        )),
930        RequestCommand::FundingRates(cmd) => DataResponse::FundingRates(FundingRatesResponse::new(
931            cmd.request_id,
932            resolve_response_client_id(cmd.client_id, used_client_id),
933            cmd.instrument_id,
934            Vec::new(),
935            Some(start),
936            Some(end),
937            ts_init,
938            cmd.params.clone(),
939        )),
940        RequestCommand::Bars(cmd) => DataResponse::Bars(BarsResponse::new(
941            cmd.request_id,
942            resolve_response_client_id(cmd.client_id, used_client_id),
943            cmd.bar_type,
944            Vec::new(),
945            Some(start),
946            Some(end),
947            ts_init,
948            cmd.params.clone(),
949        )),
950        RequestCommand::BookDeltas(cmd) => DataResponse::BookDeltas(BookDeltasResponse::new(
951            cmd.request_id,
952            resolve_response_client_id(cmd.client_id, used_client_id),
953            cmd.instrument_id,
954            Vec::new(),
955            Some(start),
956            Some(end),
957            ts_init,
958            cmd.params.clone(),
959        )),
960        RequestCommand::BookDepth(cmd) => DataResponse::BookDepth(BookDepthResponse::new(
961            cmd.request_id,
962            resolve_response_client_id(cmd.client_id, used_client_id),
963            cmd.instrument_id,
964            Vec::new(),
965            Some(start),
966            Some(end),
967            ts_init,
968            cmd.params.clone(),
969        )),
970        _ => {
971            anyhow::bail!("Cannot build empty catalog response for non-catalog-eligible request")
972        }
973    };
974
975    Ok(response)
976}
977
978fn build_quotes_catalog_response(
979    cmd: &RequestQuotes,
980    data: Vec<QuoteTick>,
981    start: UnixNanos,
982    end: UnixNanos,
983    used_client_id: Option<ClientId>,
984    ts_init: UnixNanos,
985) -> DataResponse {
986    let params = catalog_response_params(cmd.params.as_ref());
987    DataResponse::Quotes(QuotesResponse::new(
988        cmd.request_id,
989        resolve_response_client_id(cmd.client_id, used_client_id),
990        cmd.instrument_id,
991        data,
992        Some(start),
993        Some(end),
994        ts_init,
995        Some(params),
996    ))
997}
998
999fn build_trades_catalog_response(
1000    cmd: &RequestTrades,
1001    data: Vec<TradeTick>,
1002    start: UnixNanos,
1003    end: UnixNanos,
1004    used_client_id: Option<ClientId>,
1005    ts_init: UnixNanos,
1006) -> DataResponse {
1007    let params = catalog_response_params(cmd.params.as_ref());
1008    DataResponse::Trades(TradesResponse::new(
1009        cmd.request_id,
1010        resolve_response_client_id(cmd.client_id, used_client_id),
1011        cmd.instrument_id,
1012        data,
1013        Some(start),
1014        Some(end),
1015        ts_init,
1016        Some(params),
1017    ))
1018}
1019
1020fn build_funding_rates_catalog_response(
1021    cmd: &RequestFundingRates,
1022    data: Vec<FundingRateUpdate>,
1023    start: UnixNanos,
1024    end: UnixNanos,
1025    used_client_id: Option<ClientId>,
1026    ts_init: UnixNanos,
1027) -> DataResponse {
1028    let params = catalog_response_params(cmd.params.as_ref());
1029    DataResponse::FundingRates(FundingRatesResponse::new(
1030        cmd.request_id,
1031        resolve_response_client_id(cmd.client_id, used_client_id),
1032        cmd.instrument_id,
1033        data,
1034        Some(start),
1035        Some(end),
1036        ts_init,
1037        Some(params),
1038    ))
1039}
1040
1041fn build_bars_catalog_response(
1042    cmd: &RequestBars,
1043    data: Vec<Bar>,
1044    start: UnixNanos,
1045    end: UnixNanos,
1046    used_client_id: Option<ClientId>,
1047    ts_init: UnixNanos,
1048) -> DataResponse {
1049    let params = catalog_response_params(cmd.params.as_ref());
1050    DataResponse::Bars(BarsResponse::new(
1051        cmd.request_id,
1052        resolve_response_client_id(cmd.client_id, used_client_id),
1053        cmd.bar_type,
1054        data,
1055        Some(start),
1056        Some(end),
1057        ts_init,
1058        Some(params),
1059    ))
1060}
1061
1062fn build_custom_data_catalog_response(
1063    cmd: &RequestCustomData,
1064    data: Vec<CustomData>,
1065    start: UnixNanos,
1066    end: UnixNanos,
1067    ts_init: UnixNanos,
1068) -> DataResponse {
1069    let params = catalog_response_params(cmd.params.as_ref());
1070    DataResponse::Data(CustomDataResponse::new(
1071        cmd.request_id,
1072        cmd.client_id,
1073        None,
1074        cmd.data_type.clone(),
1075        data,
1076        Some(start),
1077        Some(end),
1078        ts_init,
1079        Some(params),
1080    ))
1081}
1082
1083fn build_book_deltas_catalog_response(
1084    cmd: &RequestBookDeltas,
1085    data: Vec<OrderBookDelta>,
1086    start: UnixNanos,
1087    end: UnixNanos,
1088    used_client_id: Option<ClientId>,
1089    ts_init: UnixNanos,
1090) -> DataResponse {
1091    let params = catalog_response_params(cmd.params.as_ref());
1092    DataResponse::BookDeltas(BookDeltasResponse::new(
1093        cmd.request_id,
1094        resolve_response_client_id(cmd.client_id, used_client_id),
1095        cmd.instrument_id,
1096        data,
1097        Some(start),
1098        Some(end),
1099        ts_init,
1100        Some(params),
1101    ))
1102}
1103
1104fn build_book_depth_catalog_response(
1105    cmd: &RequestBookDepth,
1106    data: Vec<OrderBookDepth10>,
1107    start: UnixNanos,
1108    end: UnixNanos,
1109    used_client_id: Option<ClientId>,
1110    ts_init: UnixNanos,
1111) -> DataResponse {
1112    let params = catalog_response_params(cmd.params.as_ref());
1113    DataResponse::BookDepth(BookDepthResponse::new(
1114        cmd.request_id,
1115        resolve_response_client_id(cmd.client_id, used_client_id),
1116        cmd.instrument_id,
1117        data,
1118        Some(start),
1119        Some(end),
1120        ts_init,
1121        Some(params),
1122    ))
1123}
1124
1125fn catalog_response_params(existing: Option<&Params>) -> Params {
1126    let mut params = existing.cloned().unwrap_or_else(Params::new);
1127    params.insert(PARAM_UPDATE_CATALOG.to_string(), Value::Bool(false));
1128    params
1129}
1130
1131fn custom_data_from_dynamic(data: Vec<Data>) -> Vec<CustomData> {
1132    data.into_iter()
1133        .filter_map(|item| match item {
1134            Data::Custom(custom) => Some(custom),
1135            other => {
1136                log::error!("Custom catalog query returned non-custom data {other:?}");
1137                None
1138            }
1139        })
1140        .collect()
1141}
1142
1143fn instrument_only_last(params: Option<&Params>) -> bool {
1144    params
1145        .and_then(|params| params.get_bool("only_last"))
1146        .unwrap_or(true)
1147}
1148
1149fn latest_instruments(data: Vec<InstrumentAny>) -> Vec<InstrumentAny> {
1150    let mut instruments: AHashMap<_, InstrumentAny> = AHashMap::new();
1151
1152    for instrument in data {
1153        let id = instrument.id();
1154        match instruments.get(&id) {
1155            Some(existing) if existing.ts_init() >= instrument.ts_init() => {}
1156            _ => {
1157                instruments.insert(id, instrument);
1158            }
1159        }
1160    }
1161
1162    let mut data: Vec<_> = instruments.into_values().collect();
1163    data.sort_by_key(|instrument| instrument.id().to_string());
1164    data
1165}
1166
1167fn instrument_response_venue(request_venue: Option<Venue>, data: &[InstrumentAny]) -> Venue {
1168    request_venue.unwrap_or_else(|| {
1169        data.iter()
1170            .map(Instrument::venue)
1171            .min_by_key(std::string::ToString::to_string)
1172            .unwrap_or_else(|| Venue::from(CATALOG_CLIENT_ID))
1173    })
1174}
1175
1176fn missing_interval_diff(start: u64, end: u64, closed_intervals: &[(u64, u64)]) -> Vec<(u64, u64)> {
1177    if closed_intervals.is_empty() {
1178        return vec![(start, end)];
1179    }
1180
1181    let mut missing = Vec::new();
1182    let mut cursor = start;
1183
1184    for &(closed_start, closed_end) in closed_intervals {
1185        if closed_end < cursor {
1186            continue;
1187        }
1188
1189        if closed_start > end {
1190            break;
1191        }
1192
1193        if closed_start > cursor {
1194            missing.push((cursor, closed_start.saturating_sub(1)));
1195        }
1196
1197        cursor = cursor.max(closed_end.saturating_add(1));
1198
1199        if cursor > end {
1200            break;
1201        }
1202    }
1203
1204    if cursor <= end {
1205        missing.push((cursor, end));
1206    }
1207
1208    missing
1209}
1210
1211fn resolve_response_client_id(
1212    request_client_id: Option<ClientId>,
1213    used_client_id: Option<ClientId>,
1214) -> ClientId {
1215    request_client_id
1216        .or(used_client_id)
1217        .unwrap_or_else(|| ClientId::new(CATALOG_CLIENT_ID))
1218}
1219
1220#[cfg(test)]
1221mod tests {
1222    use nautilus_common::messages::data::RequestJoin;
1223    use rstest::rstest;
1224
1225    use super::*;
1226
1227    #[rstest]
1228    fn test_build_empty_response_rejects_non_catalog_variant() {
1229        let request = RequestCommand::Join(RequestJoin::new(
1230            vec![UUID4::new()],
1231            None,
1232            None,
1233            UUID4::new(),
1234            UnixNanos::default(),
1235            None,
1236            None,
1237        ));
1238
1239        let result = build_empty_response(
1240            &request,
1241            UnixNanos::from(1u64),
1242            UnixNanos::from(2u64),
1243            None,
1244            UnixNanos::from(3u64),
1245        );
1246
1247        assert_eq!(
1248            result.unwrap_err().to_string(),
1249            "Cannot build empty catalog response for non-catalog-eligible request"
1250        );
1251    }
1252}