Skip to main content

apple_quant_algorithmic/backends/
databento.rs

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			// Self::Mbo => SchemaFlags::AddCancelModifyFillClear | SchemaFlags::Trades,
73			// Self::Mbp10 => SchemaFlags::LimitedAddCancelModifyFillClear | SchemaFlags::Trades,
74			// Self::Mbp1 => SchemaFlags::LimitedAddCancelModifyFillClear | SchemaFlags::Trades,
75			// Self::Tbbo => SchemaFlags::Trades,
76			Schema::Trades => SchemaFlags::Trades,
77			// Self::Ohlcv1S => SchemaFlags::OHLC1s,
78			// Self::Ohlcv1M => SchemaFlags::OHLC1m,
79			_ => 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		// Historical client only.
279		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		// Live client only.
307		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}