1use chrono::{DateTime, NaiveDateTime};
8
9use crate::models::{Bar, Timeframe};
10
11#[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 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
49pub struct BucketSpec {
50 duration_seconds: i64,
51 alignment_offset_seconds: i64,
52}
53
54impl BucketSpec {
55 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
75pub 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#[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#[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 pub fn rejected_out_of_order(&self) -> u64 {
198 self.rejected_out_of_order
199 }
200
201 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 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
259pub 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 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}