1use std::collections::BTreeMap;
2
3use chrono::{Duration, NaiveDateTime};
4
5use crate::{DataError, Result, StoredTick};
6
7#[derive(Debug, Clone, Copy, PartialEq)]
8pub struct PriorQuoteCarry {
9 pub observed_at: NaiveDateTime,
10 pub bid: f64,
11 pub ask: f64,
12 pub stale_limit: Duration,
13}
14
15#[derive(Debug, Clone, Copy, PartialEq, serde::Serialize, serde::Deserialize)]
16pub struct PriceBins {
17 pub lower: f64,
18 pub width: f64,
19 pub count: usize,
20}
21
22#[derive(Debug, Clone, PartialEq)]
23pub struct QuoteStatistics {
24 pub accepted: u64,
25 pub rejected: u64,
26 pub duplicate_provider_rows: u64,
27 pub coverage_millis: i64,
28 pub spread_last: Option<f64>,
29 pub spread_mean: Option<f64>,
30 pub spread_max: Option<f64>,
31 pub spread_p90: Option<f64>,
32 pub spread_time_mean: Option<f64>,
33 pub spread_bps_mean: Option<f64>,
34 pub quote_activity: u64,
35 pub quote_rate_per_second: Option<f64>,
36 pub mid_change_count: u64,
37 pub interarrival_mean_millis: Option<f64>,
38 pub interarrival_max_millis: Option<i64>,
39 pub interarrival_cv: Option<f64>,
40 pub direction_changes: u64,
41 pub longest_direction_streak: u64,
42 pub path_length: f64,
43 pub path_efficiency: Option<f64>,
44 pub mid_change_variance: Option<f64>,
45 pub first_high_at: Option<NaiveDateTime>,
46 pub last_high_at: Option<NaiveDateTime>,
47 pub first_low_at: Option<NaiveDateTime>,
48 pub last_low_at: Option<NaiveDateTime>,
49 pub twap: Option<f64>,
50 pub crossings: u64,
51 pub cumulative_above_millis: i64,
52 pub continuous_above_millis: i64,
53 pub elapsed_since_breakout_millis: Option<i64>,
54 pub dominant_bin: Option<usize>,
55 pub dominant_bin_center: Option<f64>,
56 pub distance_from_dominant_center: Option<f64>,
57}
58
59#[derive(Clone)]
60struct Quote {
61 ts: NaiveDateTime,
62 mid: f64,
63 spread: f64,
64}
65
66type ProviderPayload = (
67 Option<f64>,
68 Option<f64>,
69 Option<f64>,
70 Option<f64>,
71 Option<i32>,
72);
73type QuoteInterval = (NaiveDateTime, NaiveDateTime, f64, f64);
74
75#[allow(clippy::too_many_arguments)]
76pub fn aggregate_quote_statistics(
77 ticks: &[StoredTick],
78 from: NaiveDateTime,
79 to: NaiveDateTime,
80 current_price: f64,
81 captured_level: Option<f64>,
82 bins: Option<PriceBins>,
83 carry: Option<PriorQuoteCarry>,
84) -> Result<QuoteStatistics> {
85 if to <= from || !current_price.is_finite() {
86 return Err(DataError::Other(
87 "invalid quote-stat interval or current price".into(),
88 ));
89 }
90 if let Some(bins) = bins
91 && (!bins.lower.is_finite()
92 || !bins.width.is_finite()
93 || bins.width <= 0.0
94 || bins.count == 0
95 || bins.count > 4096)
96 {
97 return Err(DataError::Other("invalid price bins".into()));
98 }
99
100 let mut rows = ticks.to_vec();
101 rows.sort_by_key(|row| (row.tick.ts, row.source_ordinal));
102 let mut quotes = Vec::new();
103 let mut rejected = 0_u64;
104 let mut duplicate = 0_u64;
105 let mut identities = BTreeMap::<(String, u64), ProviderPayload>::new();
106 for row in rows {
107 row.validate()?;
108 if row.tick.ts < from || row.tick.ts >= to {
109 continue;
110 }
111 if let (Some(source), Some(sequence)) = (row.source_identity.clone(), row.provider_sequence)
112 {
113 let payload = (
114 row.tick.bid,
115 row.tick.ask,
116 row.tick.last,
117 row.tick.volume,
118 row.tick.flags,
119 );
120 if let Some(previous) = identities.get(&(source.clone(), sequence)) {
121 if previous == &payload {
122 duplicate += 1;
123 continue;
124 }
125 return Err(DataError::Other(
126 "conflicting provider-sequence quote payload".into(),
127 ));
128 }
129 identities.insert((source, sequence), payload);
130 }
131 let (Some(bid), Some(ask)) = (row.tick.bid, row.tick.ask) else {
132 rejected += 1;
133 continue;
134 };
135 if !bid.is_finite() || !ask.is_finite() || bid <= 0.0 || ask <= 0.0 || ask < bid {
136 rejected += 1;
137 continue;
138 }
139 quotes.push(Quote {
140 ts: row.tick.ts,
141 mid: (bid + ask) / 2.0,
142 spread: ask - bid,
143 });
144 }
145
146 let intervals = quote_intervals("es, from, to, carry)?;
147 let duration_ms = (to - from).num_milliseconds();
148 let coverage = intervals
149 .iter()
150 .map(|(start, end, _, _)| (*end - *start).num_milliseconds())
151 .sum::<i64>();
152 let weighted_mid = intervals
153 .iter()
154 .map(|(start, end, mid, _)| (*end - *start).num_milliseconds() as f64 * mid)
155 .sum::<f64>();
156 let weighted_spread = intervals
157 .iter()
158 .map(|(start, end, _, spread)| (*end - *start).num_milliseconds() as f64 * spread)
159 .sum::<f64>();
160 let spreads = quotes.iter().map(|quote| quote.spread).collect::<Vec<_>>();
161 let mut sorted_spreads = spreads.clone();
162 sorted_spreads.sort_by(f64::total_cmp);
163 let spread_p90 = if sorted_spreads.is_empty() {
164 None
165 } else {
166 let index = ((0.9 * sorted_spreads.len() as f64).ceil() as usize).saturating_sub(1);
167 Some(sorted_spreads[index])
168 };
169 let mids = quotes.iter().map(|quote| quote.mid).collect::<Vec<_>>();
170 let changes = mids
171 .windows(2)
172 .map(|window| window[1] - window[0])
173 .collect::<Vec<_>>();
174 let interarrival = quotes
175 .windows(2)
176 .map(|window| (window[1].ts - window[0].ts).num_milliseconds())
177 .collect::<Vec<_>>();
178 let interarrival_f64 = interarrival
179 .iter()
180 .map(|value| *value as f64)
181 .collect::<Vec<_>>();
182 let interarrival_mean = mean(&interarrival_f64);
183 let interarrival_cv = interarrival_mean.and_then(|average| {
184 if average == 0.0 {
185 None
186 } else {
187 variance(&interarrival_f64, average).map(|value| value.sqrt() / average)
188 }
189 });
190 let (direction_changes, longest_direction_streak) = direction_statistics(&changes);
191 let path_length = changes.iter().map(|value| value.abs()).sum::<f64>();
192 let path_efficiency = if path_length == 0.0 {
193 Some(0.0)
194 } else {
195 mids.first()
196 .zip(mids.last())
197 .map(|(first, last)| (last - first).abs() / path_length)
198 };
199 let (first_high_at, last_high_at, first_low_at, last_low_at) = extrema_times("es);
200 let level = level_statistics("es, &intervals, captured_level, to)?;
201 let bin = bin_statistics(&intervals, bins, current_price);
202
203 Ok(QuoteStatistics {
204 accepted: u64::try_from(quotes.len()).unwrap_or(u64::MAX),
205 rejected,
206 duplicate_provider_rows: duplicate,
207 coverage_millis: coverage,
208 spread_last: quotes.last().map(|quote| quote.spread),
209 spread_mean: mean(&spreads),
210 spread_max: spreads.iter().copied().reduce(f64::max),
211 spread_p90,
212 spread_time_mean: (coverage > 0).then_some(weighted_spread / coverage as f64),
213 spread_bps_mean: mean(
214 "es
215 .iter()
216 .map(|quote| quote.spread / quote.mid * 10_000.0)
217 .collect::<Vec<_>>(),
218 ),
219 quote_activity: u64::try_from(quotes.len()).unwrap_or(u64::MAX),
220 quote_rate_per_second: (duration_ms > 0)
221 .then_some(quotes.len() as f64 / (duration_ms as f64 / 1000.0)),
222 mid_change_count: changes.iter().filter(|value| **value != 0.0).count() as u64,
223 interarrival_mean_millis: interarrival_mean,
224 interarrival_max_millis: interarrival.iter().copied().max(),
225 interarrival_cv,
226 direction_changes,
227 longest_direction_streak,
228 path_length,
229 path_efficiency,
230 mid_change_variance: mean(&changes).and_then(|average| variance(&changes, average)),
231 first_high_at,
232 last_high_at,
233 first_low_at,
234 last_low_at,
235 twap: (coverage > 0).then_some(weighted_mid / coverage as f64),
236 crossings: level.crossings,
237 cumulative_above_millis: level.cumulative,
238 continuous_above_millis: level.continuous,
239 elapsed_since_breakout_millis: level.elapsed,
240 dominant_bin: bin.index,
241 dominant_bin_center: bin.center,
242 distance_from_dominant_center: bin.distance,
243 })
244}
245
246fn quote_intervals(
247 quotes: &[Quote],
248 from: NaiveDateTime,
249 to: NaiveDateTime,
250 carry: Option<PriorQuoteCarry>,
251) -> Result<Vec<QuoteInterval>> {
252 let mut intervals = Vec::new();
253 if let Some(carry) = carry {
254 if carry.stale_limit <= Duration::zero()
255 || !carry.bid.is_finite()
256 || !carry.ask.is_finite()
257 || carry.bid <= 0.0
258 || carry.ask < carry.bid
259 {
260 return Err(DataError::Other("invalid prior quote carry".into()));
261 }
262 let stale_end = carry
263 .observed_at
264 .checked_add_signed(carry.stale_limit)
265 .ok_or_else(|| DataError::Other("carry stale timestamp overflow".into()))?;
266 let end = quotes
267 .first()
268 .map(|quote| quote.ts)
269 .unwrap_or(to)
270 .min(to)
271 .min(stale_end);
272 if end > from && carry.observed_at <= from {
273 intervals.push((
274 from,
275 end,
276 (carry.bid + carry.ask) / 2.0,
277 carry.ask - carry.bid,
278 ));
279 }
280 }
281 for (index, quote) in quotes.iter().enumerate() {
282 let end = quotes
283 .get(index + 1)
284 .map(|next| next.ts)
285 .unwrap_or(to)
286 .min(to);
287 if end > quote.ts {
288 intervals.push((quote.ts.max(from), end, quote.mid, quote.spread));
289 }
290 }
291 Ok(intervals)
292}
293
294#[derive(Default)]
295struct LevelStatistics {
296 crossings: u64,
297 cumulative: i64,
298 continuous: i64,
299 elapsed: Option<i64>,
300}
301
302fn level_statistics(
303 quotes: &[Quote],
304 intervals: &[QuoteInterval],
305 level: Option<f64>,
306 to: NaiveDateTime,
307) -> Result<LevelStatistics> {
308 let Some(level) = level else {
309 return Ok(LevelStatistics::default());
310 };
311 if !level.is_finite() {
312 return Err(DataError::Other("captured level must be finite".into()));
313 }
314 let crossings = quotes
315 .windows(2)
316 .filter(|pair| {
317 (pair[0].mid <= level && pair[1].mid > level)
318 || (pair[0].mid >= level && pair[1].mid < level)
319 })
320 .count() as u64;
321 let mut cumulative = 0_i64;
322 let mut continuous = 0_i64;
323 let mut breakout = None;
324 let mut last_above_end = None;
325 for (start, end, mid, _) in intervals {
326 if *mid > level {
327 let duration = (*end - *start).num_milliseconds();
328 cumulative += duration;
329 breakout.get_or_insert(*start);
330 continuous = if last_above_end == Some(*start) {
331 continuous + duration
332 } else {
333 duration
334 };
335 last_above_end = Some(*end);
336 } else {
337 continuous = 0;
338 last_above_end = None;
339 }
340 }
341 Ok(LevelStatistics {
342 crossings,
343 cumulative,
344 continuous,
345 elapsed: breakout.map(|timestamp| (to - timestamp).num_milliseconds()),
346 })
347}
348
349#[derive(Default)]
350struct BinStatistics {
351 index: Option<usize>,
352 center: Option<f64>,
353 distance: Option<f64>,
354}
355
356fn bin_statistics(
357 intervals: &[QuoteInterval],
358 bins: Option<PriceBins>,
359 current_price: f64,
360) -> BinStatistics {
361 let Some(bins) = bins else {
362 return BinStatistics::default();
363 };
364 let mut dwell = vec![0_i64; bins.count];
365 for (start, end, mid, _) in intervals {
366 let relative = (*mid - bins.lower) / bins.width;
367 if relative >= 0.0 {
368 let index = relative.floor() as usize;
369 if index < bins.count {
370 dwell[index] += (*end - *start).num_milliseconds();
371 }
372 }
373 }
374 let maximum = dwell.iter().copied().max().unwrap_or(0);
375 if maximum == 0 {
376 return BinStatistics::default();
377 }
378 let index = dwell.iter().position(|value| *value == maximum).unwrap();
379 let center = bins.lower + (index as f64 + 0.5) * bins.width;
380 BinStatistics {
381 index: Some(index),
382 center: Some(center),
383 distance: Some(current_price - center),
384 }
385}
386
387fn direction_statistics(changes: &[f64]) -> (u64, u64) {
388 let signs = changes
389 .iter()
390 .map(|value| value.total_cmp(&0.0))
391 .filter(|sign| !sign.is_eq())
392 .collect::<Vec<_>>();
393 let changes = signs.windows(2).filter(|pair| pair[0] != pair[1]).count() as u64;
394 let mut streak = 0_u64;
395 let mut longest = 0_u64;
396 let mut previous = None;
397 for sign in signs {
398 if previous == Some(sign) {
399 streak += 1;
400 } else {
401 streak = 1;
402 previous = Some(sign);
403 }
404 longest = longest.max(streak);
405 }
406 (changes, longest)
407}
408
409fn mean(values: &[f64]) -> Option<f64> {
410 (!values.is_empty()).then(|| values.iter().sum::<f64>() / values.len() as f64)
411}
412
413fn variance(values: &[f64], mean: f64) -> Option<f64> {
414 (!values.is_empty()).then(|| {
415 values
416 .iter()
417 .map(|value| (value - mean).powi(2))
418 .sum::<f64>()
419 / values.len() as f64
420 })
421}
422
423fn extrema_times(
424 quotes: &[Quote],
425) -> (
426 Option<NaiveDateTime>,
427 Option<NaiveDateTime>,
428 Option<NaiveDateTime>,
429 Option<NaiveDateTime>,
430) {
431 let high = quotes.iter().map(|quote| quote.mid).reduce(f64::max);
432 let low = quotes.iter().map(|quote| quote.mid).reduce(f64::min);
433 let first_high =
434 high.and_then(|value| quotes.iter().find(|quote| quote.mid == value).map(|q| q.ts));
435 let last_high = high.and_then(|value| {
436 quotes
437 .iter()
438 .rev()
439 .find(|quote| quote.mid == value)
440 .map(|q| q.ts)
441 });
442 let first_low =
443 low.and_then(|value| quotes.iter().find(|quote| quote.mid == value).map(|q| q.ts));
444 let last_low = low.and_then(|value| {
445 quotes
446 .iter()
447 .rev()
448 .find(|quote| quote.mid == value)
449 .map(|q| q.ts)
450 });
451 (first_high, last_high, first_low, last_low)
452}