use std::range::Range;
use apple_quant_core::{log::warn, TimeDeltaExt};
use chrono::Duration;
use smallvec::SmallVec;
use tracing::instrument;
use crate::{
timestamp::{Timestamp, TradeTimestamped},
aggregation::TradeTradeTimestamp, aggregation_std::StdTrades,
binned::FlexBinnedTimePointData, binned_method::BinnedPoint, hot::FixedHotData,
};
use super::InstrumentSpec;
pub struct InstrumentData<
'instrument_data: 'aggregated_data,
'aggregated_data,
IS: InstrumentSpec,
> {
latest_trades_count: usize,
trades: Box<FixedHotData<
'instrument_data,
'aggregated_data,
100,
IS,
StdTrades<'instrument_data, IS>,
TradeTradeTimestamp<IS>,
>>,
binned_trades: FlexBinnedTimePointData<
'instrument_data,
'aggregated_data,
IS,
TradeTradeTimestamp<IS>,
StdTrades<'instrument_data, IS>,
const {
Duration::hours(1)
.num_nanoseconds()
.unwrap() as u64
},
>,
walk_previous_end_timestamp: Option<Timestamp>,
}
impl<
'instrument_data,
'aggregated_data,
IS: InstrumentSpec,
> InstrumentData<
'instrument_data,
'aggregated_data,
IS,
> {
pub fn new_trades_binned(
&mut self,
data: impl ExactSizeIterator<Item = TradeTradeTimestamp<IS>>,
) {
self.binned_trades.insert(data);
}
pub fn walk_new_trades(
&mut self,
data: impl ExactSizeIterator<Item = TradeTradeTimestamp<IS>>,
) {
self.latest_trades_count = data.len().min(self.trades.capacity());
self.trades.append_manual(data);
}
pub fn initialize_walk(
&mut self,
start_timestamp: Timestamp,
) {
self.walk_previous_end_timestamp = Some(start_timestamp);
self.trades.clear();
}
#[instrument(skip_all)]
pub async fn walk<const INTERVAL_NS: u64>(
&mut self,
walk_end_timestamp: &Timestamp,
) -> Option<&Timestamp> {
let start_timestamp = self.walk_previous_end_timestamp.unwrap();
if start_timestamp >= *walk_end_timestamp {
if start_timestamp > *walk_end_timestamp {
warn!("Walk attempted when already past end timestamp.");
}
return None;
}
let mut end_timestamp = start_timestamp +
Duration::nanoseconds_i128(INTERVAL_NS as i128).unwrap();
end_timestamp = end_timestamp.min(*walk_end_timestamp);
self.walk_previous_end_timestamp = Some(end_timestamp);
let Ok(
binned_trades,
) = self.binned_trades.iter(Range {
start: &start_timestamp,
end: &end_timestamp,
}) else {
return self.walk_previous_end_timestamp.as_ref();
};
let latest_trades: SmallVec<[TradeTradeTimestamp<IS>; 100]> = binned_trades.cloned().collect();
self.latest_trades_count = latest_trades.len().min(self.trades.capacity());
self.trades.append_manual(latest_trades.into_iter());
self.walk_previous_end_timestamp.as_ref()
}
pub fn iter_recent_latest_trades_forward(
&self,
) -> impl ExactSizeIterator<Item = &TradeTradeTimestamp<IS>> + Clone {
let count = self.latest_trades_count.min(self.trades.len());
debug_assert!(count <= self.trades.as_ref().len());
let start_idx = self.trades.len() - count;
self.trades.forward_slice()[start_idx..].into_iter()
}
pub fn iter_recent_trades_forward(
&self,
count: Option<usize>,
) -> impl ExactSizeIterator<Item = &TradeTradeTimestamp<IS>> + Clone {
let count = count
.unwrap_or(self.trades.len())
.min(self.trades.len());
debug_assert!(count <= self.trades.as_ref().len());
let start_idx = self.trades.len() - count;
self.trades.forward_slice()[start_idx..].into_iter()
}
pub fn iter_reaggregate_trade_forward<T: TradeTimestamped>(
&self,
aggregation: &T,
) -> impl Iterator<Item = &TradeTradeTimestamp<IS>> {
self.reaggregate_trade_forward_slice(aggregation).iter()
}
pub fn reaggregate_trade_forward_slice<T: TradeTimestamped>(
&self,
aggregation: &T,
) -> &[TradeTradeTimestamp<IS>] {
self.trades.forward_slice_from(|
trade_trade_timestamp,
| trade_trade_timestamp.trade_timestamp() < aggregation.trade_timestamp())
}
pub fn iter_recent_trades_backward(
&self,
count: Option<usize>,
) -> impl Iterator<Item = &TradeTradeTimestamp<IS>> {
let count = count
.unwrap_or(self.trades.len())
.min(self.trades.len());
debug_assert!(count <= self.trades.as_ref().len());
self.trades.as_ref().iter_rev().take(count)
}
pub fn trades_forward_slice_from<P: FnMut(&TradeTradeTimestamp<IS>) -> bool>(
&self,
exclude: P,
) -> &[TradeTradeTimestamp<IS>] {
self.trades.forward_slice_from(exclude)
}
pub fn trades_forward_slice(
&self,
) -> &[TradeTradeTimestamp<IS>] {
self.trades.forward_slice()
}
pub fn trades_new_recent_forward_slice(
&self,
) -> &[TradeTradeTimestamp<IS>] {
let trade_trade_timestamps = self.trades_forward_slice();
debug_assert!(self.latest_trades_count <= trade_trade_timestamps.len());
let first_idx = trade_trade_timestamps.len() - self.latest_trades_count;
&trade_trade_timestamps[first_idx..]
}
}
impl<
'instrument_data: 'aggregated_data,
'aggregated_data,
IS: InstrumentSpec,
> Default for InstrumentData<
'instrument_data,
'aggregated_data,
IS,
> {
fn default() -> Self {
Self {
walk_previous_end_timestamp: None,
trades: Box::new(FixedHotData::default()),
binned_trades: FlexBinnedTimePointData::default(),
latest_trades_count: 0,
}
}
}