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