Skip to main content

apple_quant_algorithmic/instrument/
data.rs

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