Skip to main content

data_preprocess/
quote_stats.rs

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(&quotes, 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(&quotes);
200    let level = level_statistics(&quotes, &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            &quotes
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}