Skip to main content

data_preprocess/
resample.rs

1//! Deterministic tick-to-bar aggregation shared by stored bars and in-memory replay series.
2//!
3//! Historical replay builds its analysis bars from ticks while it runs, and a research workflow needs the same bars on disk so that a parameter search over bars agrees with a confirmation run over ticks. Both paths therefore use the bucket arithmetic and accumulation rules in this module, and a parity test asserts that they produce identical bars for the same input.
4//!
5//! The module is intentionally free of storage and dataframe dependencies so that it compiles for consumers that only need the models.
6
7use chrono::{DateTime, NaiveDateTime};
8
9use crate::models::{Bar, Timeframe};
10
11/// Quote side used to derive one bar price from a tick.
12#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
13pub enum PriceBasis {
14    Bid,
15    Ask,
16    Mid,
17}
18
19impl PriceBasis {
20    pub fn parse(value: &str) -> Option<Self> {
21        match value.to_ascii_lowercase().as_str() {
22            "bid" => Some(Self::Bid),
23            "ask" => Some(Self::Ask),
24            "mid" => Some(Self::Mid),
25            _ => None,
26        }
27    }
28
29    pub fn as_str(self) -> &'static str {
30        match self {
31            Self::Bid => "bid",
32            Self::Ask => "ask",
33            Self::Mid => "mid",
34        }
35    }
36
37    /// Bar price for one two-sided quote.
38    pub fn price(self, bid: f64, ask: f64) -> f64 {
39        match self {
40            Self::Bid => bid,
41            Self::Ask => ask,
42            Self::Mid => bid + (ask - bid) / 2.0,
43        }
44    }
45}
46
47/// Fixed-duration bucket geometry.
48#[derive(Debug, Clone, Copy, PartialEq, Eq)]
49pub struct BucketSpec {
50    duration_seconds: i64,
51    alignment_offset_seconds: i64,
52}
53
54impl BucketSpec {
55    /// Build a bucket geometry, reducing the alignment offset into the bucket duration.
56    pub fn new(duration_seconds: i64, alignment_offset_seconds: i64) -> Option<Self> {
57        if duration_seconds <= 0 {
58            return None;
59        }
60        Some(Self {
61            duration_seconds,
62            alignment_offset_seconds: alignment_offset_seconds.rem_euclid(duration_seconds),
63        })
64    }
65
66    pub fn duration_seconds(self) -> i64 {
67        self.duration_seconds
68    }
69
70    pub fn alignment_offset_seconds(self) -> i64 {
71        self.alignment_offset_seconds
72    }
73}
74
75/// Half-open `[open, close)` bucket containing `timestamp`.
76///
77/// Returns `None` only when the bucket boundary would overflow the representable timestamp range.
78pub fn bucket_bounds(
79    timestamp: NaiveDateTime,
80    spec: BucketSpec,
81) -> Option<(NaiveDateTime, NaiveDateTime)> {
82    let timestamp_seconds = i128::from(timestamp.and_utc().timestamp());
83    let duration = i128::from(spec.duration_seconds);
84    let offset = i128::from(spec.alignment_offset_seconds);
85    let bucket_index = (timestamp_seconds - offset).div_euclid(duration);
86    let open_seconds = bucket_index.checked_mul(duration)?.checked_add(offset)?;
87    let close_seconds = open_seconds.checked_add(duration)?;
88    let open_seconds = i64::try_from(open_seconds).ok()?;
89    let close_seconds = i64::try_from(close_seconds).ok()?;
90    let open_time = DateTime::from_timestamp(open_seconds, 0)?.naive_utc();
91    let close_time = DateTime::from_timestamp(close_seconds, 0)?.naive_utc();
92    Some((open_time, close_time))
93}
94
95/// One bucket being accumulated.
96#[derive(Debug, Clone, PartialEq)]
97struct OpenBar {
98    open_time: NaiveDateTime,
99    close_time: NaiveDateTime,
100    open: f64,
101    high: f64,
102    low: f64,
103    close: f64,
104    tick_count: u64,
105    spread_points_sum: f64,
106    spread_samples: u64,
107}
108
109impl OpenBar {
110    fn new(open_time: NaiveDateTime, close_time: NaiveDateTime, price: f64) -> Self {
111        Self {
112            open_time,
113            close_time,
114            open: price,
115            high: price,
116            low: price,
117            close: price,
118            tick_count: 1,
119            spread_points_sum: 0.0,
120            spread_samples: 0,
121        }
122    }
123
124    fn update(&mut self, price: f64) {
125        self.high = self.high.max(price);
126        self.low = self.low.min(price);
127        self.close = price;
128        self.tick_count = self.tick_count.saturating_add(1);
129    }
130
131    fn observe_spread(&mut self, bid: f64, ask: f64, point_size: f64) {
132        if point_size <= 0.0 || !point_size.is_finite() {
133            return;
134        }
135        let spread = (ask - bid) / point_size;
136        if spread.is_finite() && spread >= 0.0 {
137            self.spread_points_sum += spread;
138            self.spread_samples += 1;
139        }
140    }
141
142    fn average_spread_points(&self) -> i32 {
143        if self.spread_samples == 0 {
144            return 0;
145        }
146        let average = self.spread_points_sum / self.spread_samples as f64;
147        if !average.is_finite() || average < 0.0 {
148            return 0;
149        }
150        average.round().min(f64::from(i32::MAX)) as i32
151    }
152}
153
154/// Accumulates ticks into fixed-duration bars using the replay engine's bucket rules.
155///
156/// A tick is accepted only when the selected price basis is available, finite, positive, and not crossed, which matches the validity rule the replay feed applies before a quote reaches the engine. An interval with no accepted tick produces no bar, so a weekend gap simply has no bars rather than empty ones.
157///
158/// Input must be chronological. A tick belonging to an already-finished bucket is dropped and counted by [`BarAggregator::rejected_out_of_order`] rather than reopening that bucket.
159///
160/// Bars record `tick_vol`, the number of accepted ticks, and leave `volume` at zero. The replay engine projects tick counts, not traded size, so a bar carries the same volume fact whether it was built here or in memory.
161#[derive(Debug, Clone)]
162pub struct BarAggregator {
163    spec: BucketSpec,
164    basis: PriceBasis,
165    exchange: String,
166    symbol: String,
167    timeframe: Timeframe,
168    point_size: f64,
169    open: Option<OpenBar>,
170    rejected_out_of_order: u64,
171}
172
173impl BarAggregator {
174    pub fn new(
175        exchange: impl Into<String>,
176        symbol: impl Into<String>,
177        timeframe: Timeframe,
178        spec: BucketSpec,
179        basis: PriceBasis,
180        point_size: f64,
181    ) -> Self {
182        Self {
183            spec,
184            basis,
185            exchange: exchange.into(),
186            symbol: symbol.into(),
187            timeframe,
188            point_size,
189            open: None,
190            rejected_out_of_order: 0,
191        }
192    }
193
194    /// Ticks dropped because they belonged to an earlier bucket than the one already open.
195    ///
196    /// A chronological source reports zero. A non-zero count means the input was not ordered, which would otherwise reopen a finished bucket and emit bars out of order.
197    pub fn rejected_out_of_order(&self) -> u64 {
198        self.rejected_out_of_order
199    }
200
201    /// Feed one tick, returning the bar that the tick just completed.
202    pub fn push(&mut self, ts: NaiveDateTime, bid: Option<f64>, ask: Option<f64>) -> Option<Bar> {
203        let (bid, ask) = (bid?, ask?);
204        if !is_executable_quote(bid, ask) {
205            return None;
206        }
207        let price = self.basis.price(bid, ask);
208        if !price.is_finite() {
209            return None;
210        }
211        let (open_time, close_time) = bucket_bounds(ts, self.spec)?;
212
213        let Some(current) = self.open.as_mut() else {
214            let mut bar = OpenBar::new(open_time, close_time, price);
215            bar.observe_spread(bid, ask, self.point_size);
216            self.open = Some(bar);
217            return None;
218        };
219        if open_time == current.open_time {
220            current.update(price);
221            current.observe_spread(bid, ask, self.point_size);
222            return None;
223        }
224        if open_time < current.open_time {
225            self.rejected_out_of_order = self.rejected_out_of_order.saturating_add(1);
226            return None;
227        }
228
229        let completed = self.open.take().expect("open bar checked above");
230        let mut next = OpenBar::new(open_time, close_time, price);
231        next.observe_spread(bid, ask, self.point_size);
232        self.open = Some(next);
233        Some(self.build(&completed))
234    }
235
236    /// Emit the bucket that is still open, for callers that deliberately want a partial bar.
237    pub fn flush(&mut self) -> Option<Bar> {
238        let open = self.open.take()?;
239        Some(self.build(&open))
240    }
241
242    fn build(&self, bar: &OpenBar) -> Bar {
243        Bar {
244            exchange: self.exchange.clone(),
245            symbol: self.symbol.clone(),
246            timeframe: self.timeframe,
247            ts: bar.open_time,
248            open: bar.open,
249            high: bar.high,
250            low: bar.low,
251            close: bar.close,
252            tick_vol: i64::try_from(bar.tick_count).unwrap_or(i64::MAX),
253            volume: 0,
254            spread: bar.average_spread_points(),
255        }
256    }
257}
258
259/// Whether a two-sided quote is executable, matching the replay feed's acceptance rule.
260pub fn is_executable_quote(bid: f64, ask: f64) -> bool {
261    bid.is_finite() && ask.is_finite() && bid > 0.0 && ask > 0.0 && ask >= bid
262}
263
264#[cfg(test)]
265mod tests {
266    #[test]
267    fn a_tick_from_a_finished_bucket_is_dropped_and_counted() {
268        let mut aggregator = hourly();
269        assert!(aggregator.push(ts(9, 0), Some(1.1), Some(1.2)).is_none());
270        let completed = aggregator.push(ts(10, 0), Some(1.3), Some(1.4)).unwrap();
271        assert_eq!(completed.ts, ts(9, 0));
272
273        // A late tick from the bucket that already closed must not reopen it.
274        assert!(aggregator.push(ts(9, 30), Some(9.9), Some(9.9)).is_none());
275        assert_eq!(aggregator.rejected_out_of_order(), 1);
276
277        let still_open = aggregator.flush().unwrap();
278        assert_eq!(still_open.ts, ts(10, 0));
279        assert_eq!(still_open.high, 1.3);
280        assert_eq!(still_open.tick_vol, 1);
281    }
282
283    #[test]
284    fn bars_report_tick_counts_and_leave_traded_volume_at_zero() {
285        let mut aggregator = hourly();
286        aggregator.push(ts(9, 0), Some(1.1), Some(1.2));
287        aggregator.push(ts(9, 30), Some(1.15), Some(1.25));
288        let bar = aggregator.flush().unwrap();
289        assert_eq!(bar.tick_vol, 2);
290        assert_eq!(bar.volume, 0);
291    }
292
293    use super::*;
294    use chrono::NaiveDate;
295
296    fn ts(hour: u32, minute: u32) -> NaiveDateTime {
297        NaiveDate::from_ymd_opt(2026, 6, 1)
298            .unwrap()
299            .and_hms_opt(hour, minute, 0)
300            .unwrap()
301    }
302
303    fn hourly() -> BarAggregator {
304        BarAggregator::new(
305            "demo",
306            "EURUSD",
307            Timeframe::H1,
308            BucketSpec::new(3600, 0).unwrap(),
309            PriceBasis::Bid,
310            1.0e-5,
311        )
312    }
313
314    #[test]
315    fn bucket_bounds_align_to_the_offset() {
316        let spec = BucketSpec::new(86_400, 79_200).unwrap();
317        let (open, close) = bucket_bounds(ts(23, 0), spec).unwrap();
318        assert_eq!(open, ts(22, 0));
319        assert_eq!(close, ts(22, 0) + chrono::Duration::days(1));
320    }
321
322    #[test]
323    fn an_offset_is_reduced_into_the_duration() {
324        let spec = BucketSpec::new(3600, 7200).unwrap();
325        assert_eq!(spec.alignment_offset_seconds(), 0);
326    }
327
328    #[test]
329    fn a_completed_bucket_is_emitted_when_the_next_one_opens() {
330        let mut aggregator = hourly();
331        assert!(
332            aggregator
333                .push(ts(10, 0), Some(1.1), Some(1.10002))
334                .is_none()
335        );
336        assert!(
337            aggregator
338                .push(ts(10, 30), Some(1.2), Some(1.20002))
339                .is_none()
340        );
341        let bar = aggregator
342            .push(ts(11, 0), Some(1.05), Some(1.05002))
343            .unwrap();
344        assert_eq!(bar.ts, ts(10, 0));
345        assert_eq!(bar.open, 1.1);
346        assert_eq!(bar.high, 1.2);
347        assert_eq!(bar.low, 1.1);
348        assert_eq!(bar.close, 1.2);
349        assert_eq!(bar.tick_vol, 2);
350    }
351
352    #[test]
353    fn an_empty_interval_produces_no_bar() {
354        let mut aggregator = hourly();
355        aggregator.push(ts(10, 0), Some(1.1), Some(1.10002));
356        let bar = aggregator
357            .push(ts(13, 0), Some(1.2), Some(1.20002))
358            .unwrap();
359        assert_eq!(bar.ts, ts(10, 0));
360        assert!(aggregator.flush().unwrap().ts == ts(13, 0));
361    }
362
363    #[test]
364    fn invalid_and_one_sided_ticks_are_skipped() {
365        let mut aggregator = hourly();
366        assert!(aggregator.push(ts(10, 0), None, Some(1.1)).is_none());
367        assert!(aggregator.push(ts(10, 1), Some(1.1), None).is_none());
368        assert!(aggregator.push(ts(10, 2), Some(1.2), Some(1.1)).is_none());
369        assert!(aggregator.push(ts(10, 3), Some(-1.0), Some(1.1)).is_none());
370        assert!(aggregator.flush().is_none());
371    }
372
373    #[test]
374    fn the_average_spread_is_stored_in_points() {
375        let mut aggregator = hourly();
376        aggregator.push(ts(10, 0), Some(1.10000), Some(1.10001));
377        aggregator.push(ts(10, 1), Some(1.10000), Some(1.10003));
378        let bar = aggregator
379            .push(ts(11, 0), Some(1.10000), Some(1.10001))
380            .unwrap();
381        assert_eq!(bar.spread, 2);
382    }
383
384    #[test]
385    fn the_mid_basis_averages_both_sides() {
386        let mut aggregator = BarAggregator::new(
387            "demo",
388            "EURUSD",
389            Timeframe::H1,
390            BucketSpec::new(3600, 0).unwrap(),
391            PriceBasis::Mid,
392            1.0e-5,
393        );
394        aggregator.push(ts(10, 0), Some(1.0), Some(1.2));
395        let bar = aggregator.push(ts(11, 0), Some(1.0), Some(1.0)).unwrap();
396        assert_eq!(bar.open, 1.1);
397    }
398}