Skip to main content

apple_quant_algorithmic/instrument/
data.rs

1use std::range::Range;
2
3use apple_quant_core::{log::warn, TimeDeltaExt};
4
5use chrono::Duration;
6use smallvec::SmallVec;
7use tracing::instrument;
8
9use crate::{
10	timestamp::{Timestamp, TradeTimestamped},
11	aggregation::TradeTradeTimestamp, aggregation_std::StdTrades,
12	binned::FlexBinnedTimePointData, binned_method::BinnedPoint, hot::FixedHotData,
13};
14
15use super::InstrumentSpec;
16
17/// Primary hot storage of securities enabled foundational data used for further aggregation and
18/// permanent storage.
19pub struct InstrumentData<
20	'instrument_data: 'aggregated_data,
21	'aggregated_data,
22	IS: InstrumentSpec,
23> {
24	latest_trades_count: usize,
25
26	trades: Box<FixedHotData<
27		'instrument_data,
28		'aggregated_data,
29		100,
30		IS,
31		StdTrades<'instrument_data, IS>,
32		TradeTradeTimestamp<IS>,
33	>>,
34
35	binned_trades: FlexBinnedTimePointData<
36		'instrument_data,
37		'aggregated_data,
38		IS,
39		TradeTradeTimestamp<IS>,
40		StdTrades<'instrument_data, IS>,
41		const {
42			Duration::hours(1)
43				.num_nanoseconds()
44				.unwrap() as u64
45		},
46	>,
47
48	walk_previous_end_timestamp: Option<Timestamp>,
49}
50
51impl<
52	'instrument_data,
53	'aggregated_data,
54	IS: InstrumentSpec,
55> InstrumentData<
56	'instrument_data,
57	'aggregated_data,
58	IS,
59> {
60	pub fn new_trades_binned(
61		&mut self,
62		data: impl ExactSizeIterator<Item = TradeTradeTimestamp<IS>>,
63	) {
64		self.binned_trades.insert(data);
65	}
66
67	pub fn walk_new_trades(
68		&mut self,
69		data: impl ExactSizeIterator<Item = TradeTradeTimestamp<IS>>,
70	) {
71		self.latest_trades_count = data.len().min(self.trades.capacity());
72		self.trades.append_manual(data);
73	}
74
75	pub fn initialize_walk(
76		&mut self,
77		start_timestamp: Timestamp,
78	) {
79		self.walk_previous_end_timestamp = Some(start_timestamp);
80		self.trades.clear();
81	}
82
83	#[instrument(skip_all)]
84	pub async fn walk<const INTERVAL_NS: u64>(
85		&mut self,
86		walk_end_timestamp: &Timestamp,
87	) -> Option<&Timestamp> {
88		let start_timestamp = self.walk_previous_end_timestamp.unwrap();
89
90		if start_timestamp >= *walk_end_timestamp {
91			if start_timestamp > *walk_end_timestamp {
92				warn!("Walk attempted when already past end timestamp.");
93			}
94
95			return None;
96		}
97
98		let mut end_timestamp = start_timestamp +
99			Duration::nanoseconds_i128(INTERVAL_NS as i128).unwrap();
100
101		end_timestamp = end_timestamp.min(*walk_end_timestamp);
102		self.walk_previous_end_timestamp = Some(end_timestamp);
103
104		let Ok(
105			binned_trades,
106		) = self.binned_trades.iter(Range {
107			start: &start_timestamp,
108			end: &end_timestamp,
109		}) else {
110			return self.walk_previous_end_timestamp.as_ref();
111		};
112
113		let latest_trades: SmallVec<[TradeTradeTimestamp<IS>; 100]> = binned_trades.cloned().collect();
114
115		self.latest_trades_count = latest_trades.len().min(self.trades.capacity());
116		self.trades.append_manual(latest_trades.into_iter());
117		self.walk_previous_end_timestamp.as_ref()
118	}
119
120	pub fn iter_recent_latest_trades_forward(
121		&self,
122	) -> impl ExactSizeIterator<Item = &TradeTradeTimestamp<IS>> + Clone {
123		let count = self.latest_trades_count.min(self.trades.len());
124
125		debug_assert!(count <= self.trades.as_ref().len());
126
127		let start_idx = self.trades.len() - count;
128
129		self.trades.forward_slice()[start_idx..].into_iter()
130	}
131
132	pub fn iter_recent_trades_forward(
133		&self,
134		count: Option<usize>,
135	) -> impl ExactSizeIterator<Item = &TradeTradeTimestamp<IS>> + Clone {
136		let count = count
137			.unwrap_or(self.trades.len())
138			.min(self.trades.len());
139
140		debug_assert!(count <= self.trades.as_ref().len());
141
142		let start_idx = self.trades.len() - count;
143
144		self.trades.forward_slice()[start_idx..].into_iter()
145	}
146
147	pub fn iter_reaggregate_trade_forward<T: TradeTimestamped>(
148		&self,
149		aggregation: &T,
150	) -> impl Iterator<Item = &TradeTradeTimestamp<IS>> {
151		self.reaggregate_trade_forward_slice(aggregation).iter()
152	}
153
154	pub fn reaggregate_trade_forward_slice<T: TradeTimestamped>(
155		&self,
156		aggregation: &T,
157	) -> &[TradeTradeTimestamp<IS>] {
158		self.trades.forward_slice_from(|
159			trade_trade_timestamp,
160		| trade_trade_timestamp.trade_timestamp() < aggregation.trade_timestamp())
161	}
162
163	pub fn iter_recent_trades_backward(
164		&self,
165		count: Option<usize>,
166	) -> impl Iterator<Item = &TradeTradeTimestamp<IS>> {
167		let count = count
168			.unwrap_or(self.trades.len())
169			.min(self.trades.len());
170
171		debug_assert!(count <= self.trades.as_ref().len());
172		self.trades.as_ref().iter_rev().take(count)
173	}
174
175	pub fn trades_forward_slice_from<P: FnMut(&TradeTradeTimestamp<IS>) -> bool>(
176		&self,
177		exclude: P,
178	) -> &[TradeTradeTimestamp<IS>] {
179		self.trades.forward_slice_from(exclude)
180	}
181
182	pub fn trades_forward_slice(
183		&self,
184	) -> &[TradeTradeTimestamp<IS>] {
185		self.trades.forward_slice()
186	}
187
188	pub fn trades_new_recent_forward_slice(
189		&self,
190	) -> &[TradeTradeTimestamp<IS>] {
191		let trade_trade_timestamps = self.trades_forward_slice();
192
193		debug_assert!(self.latest_trades_count <= trade_trade_timestamps.len());
194
195		let first_idx = trade_trade_timestamps.len() - self.latest_trades_count;
196
197		&trade_trade_timestamps[first_idx..]
198	}
199}
200
201impl<
202	'instrument_data: 'aggregated_data,
203	'aggregated_data,
204	IS: InstrumentSpec,
205> Default for InstrumentData<
206	'instrument_data,
207	'aggregated_data,
208	IS,
209> {
210	fn default() -> Self {
211		Self {
212			walk_previous_end_timestamp: None,
213			trades: Box::new(FixedHotData::default()),
214			binned_trades: FlexBinnedTimePointData::default(),
215			latest_trades_count: 0,
216		}
217	}
218}