1mod decoder;
2mod encoder;
3mod historical;
4mod live;
5mod symbology;
6
7use std::{marker::PhantomData, path::PathBuf, range::Range};
8
9#[allow(unused_imports)]
10use apple_quant_core::log::trace;
11
12use apple_quant_core::{
13 log::{error, info},
14 UnwrapResultExt,
15};
16
17use databento::{
18 dbn::{
19 decode::AsyncDbnDecoder, RecordRef, RecordRefEnum, Schema, Side,
20 FIXED_PRICE_SCALE,
21 },
22 DateTimeLike,
23};
24
25use smallstr::SmallString;
26
27use time::{Duration, OffsetDateTime, Time, UtcDateTime, UtcOffset};
28
29use tokio::{
30 io::{BufReader, BufWriter},
31 fs::File, task::yield_now, time::sleep,
32};
33
34use crate::{
35 aggregation::{Trade, TradeTradeTimestamp},
36 backend::{
37 DataBackend, HistoricalDataBackend, MarketDataDecoder, MarketDataDecoderProvider,
38 OrderIdGenerator, OrdersBackend, OrdersBackendUpdate, OrdersBackendUpdateRecycle,
39 RealtimeDataBackend,
40 },
41 instrument::{
42 InstrumentData, InstrumentSpec, InstrumentTicker, IsFloatingPoint, XSpec,
43 },
44 order::{DeferredOrderActions, StateGoal},
45 order_manager::{OrderManager, OrdersCapacitySpec},
46 points::{
47 Subpoints, SubpointsType, VolumeFromWholePoints, WholePoints, WholePointsType,
48 },
49 timestamp::{
50 TickTimestamp, Timestamp, TimestampRangeExcluded, TimestampRangeExclusive,
51 Timestamped, TradeTimestamp, TradeTimestamped, UtcNs,
52 },
53 volume::{
54 AggressiveVolume, AggressorSide, DirectionalExposure, DirectionlessVolume,
55 VolumeConstruct,
56 },
57 iso_string_date, price::AbsolutePrice, schema::SchemaFlags, strategy::Strategy,
58 Frontend,
59};
60
61use decoder::*;
62use encoder::*;
63use historical::*;
64use live::*;
65use symbology::*;
66
67impl From<&Schema> for SchemaFlags {
68 fn from(
69 value: &Schema,
70 ) -> Self {
71 match value {
72 Schema::Trades => SchemaFlags::Trades,
77 _ => SchemaFlags::empty(),
80 }
81 }
82}
83
84impl From<Schema> for SchemaFlags {
85 fn from(
86 value: Schema,
87 ) -> Self {
88 Self::from(&value)
89 }
90}
91
92pub(crate) fn schema_value(
93 schema: &Schema,
94) -> u8 {
95 match schema {
96 Schema::Mbo => 8,
97 Schema::Mbp10 => 7,
98 Schema::Mbp1 => 6,
99 Schema::Tbbo => 5,
100 Schema::Tcbbo => 4,
101 Schema::Trades => 3,
102 Schema::Ohlcv1S => 2,
103 Schema::Ohlcv1M => 1,
104 _ => 0,
105 }
106}
107
108impl DateTimeLike for Timestamp {
109 fn to_date_time(
110 self,
111 ) -> OffsetDateTime {
112 OffsetDateTime::from_unix_timestamp_nanos(
113 *self.as_timestamp_type(),
114 ).unwrap()
115 }
116}
117
118pub struct Databento<
119 'instrument_data,
120 'aggregated_data,
121 IS: InstrumentSpec,
122 OB: OrdersBackend<IS>,
123 OrdersCS: OrdersCapacitySpec,
124 S: Strategy<IS, OB, OrdersCS>,
125> {
126 frontend: Frontend<'instrument_data, 'aggregated_data, IS, OB, OrdersCS, S>,
127
128 orders_backend_update_receiver: thingbuf::mpsc::Receiver<
129 OrdersBackendUpdate<IS>,
130 OrdersBackendUpdateRecycle,
131 >,
132
133 key: Option<SmallString<[u8; 64]>>,
134 realtime: Option<DatabentoLive>,
135 historical: Option<DatabentoHistorical>,
136}
137
138impl<
139 'instrument_data,
140 'aggregated_data,
141 IS: InstrumentSpec,
142 OB: OrdersBackend<IS>,
143 OrdersCS: OrdersCapacitySpec,
144 S: Strategy<IS, OB, OrdersCS>,
145> Databento<
146 'instrument_data,
147 'aggregated_data,
148 IS,
149 OB,
150 OrdersCS,
151 S,
152> {
153 const HISTORICAL_LIVE_OVERLAP_SHORT: Duration = Duration::seconds(5);
154 const HISTORICAL_LIVE_OVERLAP_LONG: Duration = Duration::minutes(10);
155
156 pub fn new(
157 key: Option<&str>,
158 frontend: Frontend<'instrument_data, 'aggregated_data, IS, OB, OrdersCS, S>,
159 orders_backend_update_receiver: thingbuf::mpsc::Receiver<
160 OrdersBackendUpdate<IS>,
161 OrdersBackendUpdateRecycle,
162 >,
163 ) -> Self {
164 Databento::<'instrument_data, 'aggregated_data> {
165 frontend,
166 orders_backend_update_receiver,
167 key: key.map(|key| SmallString::from_str(key)),
168 realtime: None,
169 historical: None,
170 }
171 }
172
173 async fn new_realtime(
174 &mut self,
175 instrument_ticker: &InstrumentTicker,
176 start_timestamp: impl DateTimeLike,
177 ) -> Result<&mut DatabentoLive, ()> {
178 if let Some(
179 mut realtime_databento,
180 ) = self.realtime.take() {
181 let _ = realtime_databento.client.close().await;
182 }
183
184 let Some(
185 key,
186 ) = &self.key else {
187 return Err(());
188 };
189
190 let Ok((
191 realtime_databento,
192 _schema_flags,
193 )) = DatabentoLive::new(
194 instrument_ticker,
195 start_timestamp,
196 &[Schema::Trades],
197 key.as_str(),
198 ).await else {
199 return Err(());
200 };
201
202 unsafe {
203 self.realtime = Some(realtime_databento);
204 Ok(self.realtime.as_mut().unwrap_unchecked())
205 }
206 }
207
208 async fn current_realtime(
209 &mut self,
210 ) -> Option<&mut DatabentoLive> {
211 self.realtime.as_mut()
212 }
213
214 async fn initialize_historical(
215 &mut self,
216 ) -> Result<&mut DatabentoHistorical, ()> {
217 let Some(
218 key,
219 ) = &self.key else {
220 return Err(());
221 };
222
223 self.historical = Some(DatabentoHistorical::new(key.as_str()));
224
225 unsafe {
226 Ok(self.historical.as_mut().unwrap_unchecked())
227 }
228 }
229
230 async fn databento_historical(
231 &mut self,
232 ) -> Result<&mut DatabentoHistorical, ()> {
233 if self.historical.is_none() {
234 return self.initialize_historical().await;
235 }
236
237 let Some(
238 databento_historical,
239 ) = self.historical.as_mut() else {
240 return Err(());
241 };
242
243 Ok(databento_historical)
244 }
245}
246
247impl<
248 'instrument_data,
249 'aggregated_data,
250 IS: InstrumentSpec,
251 OB: OrdersBackend<IS>,
252 OrdersCS: OrdersCapacitySpec,
253 S: Strategy<IS, OB, OrdersCS>,
254> HistoricalDataBackend<IS> for Databento<
255 'instrument_data,
256 'aggregated_data,
257 IS,
258 OB,
259 OrdersCS,
260 S,
261> {
262 async fn fetch_once(
263 &mut self,
264 instrument_ticker: &InstrumentTicker,
265 utc_date_time_range: Range<UtcDateTime>,
266 ) -> impl IntoIterator<Item = TradeTradeTimestamp<IS>> {
267 debug_assert!(utc_date_time_range.start <= utc_date_time_range.end);
268
269 let now = UtcDateTime::now();
270 let exact_cutoff = now - Duration::hours(24);
271 let live_safe_cutoff = exact_cutoff + Self::HISTORICAL_LIVE_OVERLAP_SHORT;
272 let historical_safe_cutoff = exact_cutoff - Self::HISTORICAL_LIVE_OVERLAP_SHORT;
273
274 debug_assert!(utc_date_time_range.end < (now + Duration::seconds(5)));
275
276 let mut market_data = Vec::with_capacity(1_000_000);
277
278 if (
280 utc_date_time_range.start <= historical_safe_cutoff &&
281 utc_date_time_range.end <= historical_safe_cutoff
282 ) {
283 let historical_databento = self.databento_historical().await.unwrap();
284
285 let mut stream = historical_databento
286 .stream(
287 instrument_ticker,
288 utc_date_time_range,
289 ).await.unwrap();
290
291 while let Ok(Some(
292 record_ref,
293 )) = stream.decode_record_ref().await {
294 let Some(
295 trade_trade_timestamp,
296 ) = process_record_trades(record_ref) else {
297 continue;
298 };
299
300 market_data.push(trade_trade_timestamp);
301 }
302
303 return market_data;
304 }
305
306 if (
308 utc_date_time_range.start >= live_safe_cutoff &&
309 utc_date_time_range.end >= live_safe_cutoff
310 ) {
311 let realtime_databento = self.new_realtime(
312 instrument_ticker,
313 utc_date_time_range.start.to_offset(UtcOffset::UTC),
314 ).await.unwrap();
315
316 while let Ok(Some(
317 record_ref,
318 )) = realtime_databento.client.next_record().await {
319 let Some(
320 trade_trade_timestamp,
321 ) = process_record_trades(record_ref) else {
322 continue;
323 };
324
325 let trade_date_time = trade_trade_timestamp
326 .trade_timestamp()
327 .to_date_time().to_utc();
328
329 if trade_date_time >= utc_date_time_range.end {
330 break;
331 }
332
333 market_data.push(trade_trade_timestamp);
334 }
335
336 let _ = realtime_databento.client.close().await;
337 self.realtime = None;
338 return market_data;
339 }
340
341 let historical_start = if utc_date_time_range.start < historical_safe_cutoff {
342 utc_date_time_range.start
343 } else {
344 historical_safe_cutoff
345 };
346
347 let realtime_end = if utc_date_time_range.end > live_safe_cutoff {
348 utc_date_time_range.end
349 } else {
350 live_safe_cutoff
351 };
352
353 let realtime_databento = self.new_realtime(
354 instrument_ticker,
355 OffsetDateTime::UNIX_EPOCH,
356 ).await.unwrap();
357
358 let mut realtime_satisfied_start = false;
359
360 while let Ok(Some(
361 record_ref,
362 )) = realtime_databento.client.next_record().await {
363 let Some(
364 trade_trade_timestamp,
365 ) = process_record_trades(record_ref) else {
366 continue;
367 };
368
369 let trade_date_time = trade_trade_timestamp
370 .trade_timestamp()
371 .to_date_time().to_utc();
372
373 if trade_date_time >= realtime_end {
374 break;
375 }
376
377 if trade_date_time < historical_start {
378 realtime_satisfied_start = true;
379 continue;
380 }
381
382 market_data.push(trade_trade_timestamp);
383 }
384
385 let _ = realtime_databento.client.close().await;
386 self.realtime = None;
387
388 if realtime_satisfied_start {
389 return market_data;
390 }
391
392 let first_live_timestamp = market_data
393 .first()
394 .unwrap()
395 .trade_timestamp()
396 .to_date_time().to_utc();
397
398 debug_assert!(first_live_timestamp < live_safe_cutoff);
399
400 let historical_databento = self.databento_historical().await.unwrap();
401 let start_utc = UtcDateTime::now();
402
403 while UtcDateTime::now() - start_utc < Duration::seconds(30) {
404 let available_end = historical_databento.available_end().await.unwrap();
405
406 if available_end > first_live_timestamp {
407 break;
408 }
409
410 sleep(std::time::Duration::from_millis(100)).await;
411 }
412
413 let mut historical_market_data = Vec::with_capacity(10_000);
414
415 let mut stream = historical_databento
416 .stream(
417 instrument_ticker,
418 Range {
419 start: historical_start,
420 end: (first_live_timestamp + Duration::nanoseconds(1)),
421 },
422 ).await.unwrap();
423
424 while let Ok(Some(
425 record_ref,
426 )) = stream.decode_record_ref().await {
427 let Some(
428 trade_trade_timestamp,
429 ) = process_record_trades(record_ref) else {
430 continue;
431 };
432
433 historical_market_data.push(trade_trade_timestamp);
434 }
435
436 let inclusive_start_live_idx = market_data.iter().position(|
437 trade_trade_timestamp,
438 | {
439 trade_trade_timestamp
440 .trade_timestamp()
441 .to_date_time().to_utc() > first_live_timestamp
442 }).unwrap();
443
444 historical_market_data.extend(
445 market_data.into_iter().skip(inclusive_start_live_idx),
446 );
447
448 historical_market_data
449 }
450
451 async fn fetch_save_day(
452 &mut self,
453 instrument_ticker: &InstrumentTicker,
454 date: time::Date,
455 ) -> impl IntoIterator<Item = TradeTradeTimestamp<IS>> {
456 let start = UtcDateTime::new(date, Time::MIDNIGHT);
457 let end = UtcDateTime::new(date.next_day().unwrap(), Time::MIDNIGHT);
458
459 let utc_date_time_range = Range { start, end };
460
461 debug_assert!(utc_date_time_range.start <= utc_date_time_range.end);
462
463 let mut iso_string_date = iso_string_date(date).unwrap();
464
465 iso_string_date.push_str(".dbn");
466
467 let mut path_buf = PathBuf::new();
468
469 path_buf.push("market-data");
470 path_buf.push(IS::INSTRUMENT_KIND.as_str());
471 path_buf.push(IS::BARE_SYMBOL);
472 path_buf.push("trades");
473 path_buf.push(iso_string_date);
474
475 let first_buf_writer = BufWriter::new(File::create(path_buf).await.unwrap());
476
477 let mut encoder = DatabentoMarketDataEncoder::new(
478 first_buf_writer,
479 &utc_date_time_range,
480 instrument_ticker,
481 ).await;
482
483 let now = UtcDateTime::now();
484 let exact_cutoff = now - Duration::hours(24);
485 let live_safe_cutoff = exact_cutoff + Self::HISTORICAL_LIVE_OVERLAP_SHORT;
486 let historical_safe_cutoff = exact_cutoff - Self::HISTORICAL_LIVE_OVERLAP_SHORT;
487
488 debug_assert!(utc_date_time_range.end < (now + Duration::seconds(5)));
489
490 let mut market_data = Vec::with_capacity(1_000_000);
491
492 if (
493 utc_date_time_range.start > historical_safe_cutoff ||
494 utc_date_time_range.end > historical_safe_cutoff
495 ) {
496 unimplemented!();
497 }
498
499 let historical_databento = self.databento_historical().await.unwrap();
500
501 let mut stream = historical_databento
502 .stream(
503 instrument_ticker,
504 utc_date_time_range,
505 ).await.unwrap();
506
507 while let Ok(Some(
508 record_ref,
509 )) = stream.decode_record_ref().await {
510 let Some(
511 trade_trade_timestamp,
512 ) = process_record_trades(record_ref) else {
513 continue;
514 };
515
516 encoder.encode_record_ref(record_ref).await;
517 market_data.push(trade_trade_timestamp);
518 }
519
520 encoder.shutdown().await;
521 market_data
522 }
523}
524
525impl<
526 'instrument_data,
527 'aggregated_data,
528 IS: InstrumentSpec,
529 OB: OrdersBackend<IS>,
530 OrdersCS: OrdersCapacitySpec,
531 S: Strategy<IS, OB, OrdersCS> + Send,
532> RealtimeDataBackend<
533 'instrument_data,
534 'aggregated_data,
535 IS,
536 OB,
537 OrdersCS,
538 S,
539> for Databento<
540 'instrument_data,
541 'aggregated_data,
542 IS,
543 OB,
544 OrdersCS,
545 S,
546> {
547 async fn initialize_realtime(
548 &mut self,
549 instrument_ticker: &InstrumentTicker,
550 recent_trade_timestamp: Option<&TradeTimestamp>,
551 ) {
552 let start_timestamp = recent_trade_timestamp
553 .map(|
554 trade_timestamp,
555 | trade_timestamp.to_date_time())
556 .unwrap_or_else(|| OffsetDateTime::now_utc());
557
558 let _ = self.new_realtime(
559 instrument_ticker,
560 start_timestamp,
561 ).await;
562 }
563
564 async fn poll_realtime(
565 &mut self,
566 instrument_data: &'instrument_data mut InstrumentData<
567 'instrument_data,
568 'aggregated_data,
569 IS,
570 >,
571 strategy: &mut S,
572 order_manager: &OrderManager<IS, OrdersCS>,
573 directional_exposure: &DirectionalExposure<IS>,
574 state_goal: &mut StateGoal,
575 deferred_order_actions: &mut DeferredOrderActions<IS>,
576 order_id_generator: &mut OrderIdGenerator,
577 recent_trade_timestamp: Option<&TradeTimestamp>,
578 ) {
579 let start_recent_trade_timestamp = recent_trade_timestamp.cloned();
580 let mut last_recent_trade_timestamp = recent_trade_timestamp.cloned();
581
582 #[rustfmt::skip]
583 let Some(
584 realtime_databento,
585 ) = &mut self.realtime else {
586 return;
587 };
588
589 loop {
590 let next_record = realtime_databento.client.try_next_record();
591
592 let record_ref = match next_record {
593 Err(
594 error,
595 ) => {
596 error!("{error}");
597 break;
598 },
599 Ok(
600 record_ref,
601 ) => record_ref,
602 };
603
604 if let Some(
605 record_ref,
606 ) = record_ref {
607 let Some(
608 trade_trade_timestamp,
609 ) = process_record_trades::<IS>(record_ref) else {
610 continue;
611 };
612
613 last_recent_trade_timestamp = Some(
614 trade_trade_timestamp.trade_timestamp(),
615 );
616
617 instrument_data.new_trades_binned([trade_trade_timestamp].into_iter());
618 continue;
619 }
620
621 let size = match realtime_databento.client.fill_buf().await {
622 Err(
623 error,
624 ) => {
625 error!("{error}");
626 break;
627 },
628 Ok(size) => size,
629 };
630
631 if size == 0 {
632 break;
633 }
634
635 if realtime_databento.client.is_closed() {
636 unimplemented!();
637 }
638 }
639
640 let Some(
641 last_recent_trade_timestamp,
642 ) = last_recent_trade_timestamp else {
643 return;
644 };
645
646 let start_recent_trade_timestamp = start_recent_trade_timestamp.unwrap_or(
647 last_recent_trade_timestamp,
648 );
649
650 let tick_timestamp = TickTimestamp::from_timestamp(
651 start_recent_trade_timestamp.timestamp(),
652 );
653
654 strategy.data_update(
655 &tick_timestamp,
656 instrument_data,
657 order_manager,
658 directional_exposure,
659 state_goal,
660 deferred_order_actions,
661 order_id_generator,
662 ).await;
663 }
664}
665
666impl<
667 'instrument_data,
668 'aggregated_data,
669 IS: InstrumentSpec,
670 OB: OrdersBackend<IS> + Send,
671 OrdersCS: OrdersCapacitySpec,
672 S: Strategy<IS, OB, OrdersCS> + Send,
673> DataBackend<
674 'instrument_data,
675 'aggregated_data,
676 IS,
677 OB,
678 OrdersCS,
679 S,
680> for Databento<
681 'instrument_data,
682 'aggregated_data,
683 IS,
684 OB,
685 OrdersCS,
686 S,
687> {
688 fn new(
689 data_key: Option<&str>,
690 frontend: Frontend<'instrument_data, 'aggregated_data, IS, OB, OrdersCS, S>,
691 orders_backend_update_receiver: thingbuf::mpsc::Receiver<
692 OrdersBackendUpdate<IS>,
693 OrdersBackendUpdateRecycle,
694 >,
695 ) -> Self {
696 Databento::new(data_key, frontend, orders_backend_update_receiver)
697 }
698
699 fn file_name_postpend() -> &'static str {
700 ".dbn"
701 }
702
703 async fn backtest(
704 &mut self,
705 walk_range: TimestampRangeExclusive,
706 ) {
707 let mut tick_timestamp = TickTimestamp::from_timestamp(
708 walk_range.start.as_timestamp(),
709 );
710
711 self.frontend.initialize(&tick_timestamp).await;
712
713 loop {
714 yield_now().await;
715
716 unsafe {
717 self.frontend
718 .tick(
719 &mut self.orders_backend_update_receiver,
720 &mut tick_timestamp,
721 ).await.unwrap_unchecked_release()
722 };
723 }
724
725 #[cfg(feature = "log-trace-frontend")]
726 trace!("Walk loop ended.");
727 }
728
729 async fn decode_market_data(
730 buf_readers: impl ExactSizeIterator<Item = BufReader<File>>,
731 ) -> Self::MarketDataDecoder {
732 DatabentoMarketDataDecoder::new(buf_readers).await.unwrap()
733 }
734
735 fn historical_mut(
736 &mut self,
737 ) -> Option<Result<&mut impl HistoricalDataBackend<IS>, ()>> {
738 Some(Ok(self))
739 }
740}
741
742impl<
743 'instrument_data,
744 'aggregated_data,
745 IS: InstrumentSpec,
746 OB: OrdersBackend<IS>,
747 OrdersCS: OrdersCapacitySpec,
748 S: Strategy<IS, OB, OrdersCS>,
749> MarketDataDecoderProvider<IS> for Databento<
750 'instrument_data,
751 'aggregated_data,
752 IS,
753 OB,
754 OrdersCS,
755 S,
756> {
757 type MarketDataDecoder = DatabentoMarketDataDecoder<IS>;
758}