apple_quant_algorithmic/instrument/
data.rs1use 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
17pub 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}