1use 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 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 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 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 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 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 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 _ => 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}