1use crate::model::Bar;
2use std::fmt;
3
4#[cfg(feature = "serde")]
5use serde::{Deserialize, Serialize};
6
7const SECONDS_PER_DAY: i64 = 86_400;
8
9#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
11#[cfg_attr(feature = "serde", derive(Serialize, Deserialize))]
12pub enum Timeframe {
13 Second(u32),
14 Minute(u32),
15 Hour(u32),
16 Day(u32),
17 Week(u32),
18 Month(u32),
19}
20
21#[derive(Debug, Clone, Copy, PartialEq, Eq)]
23pub struct TimeframeError;
24
25impl fmt::Display for TimeframeError {
26 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
27 f.write_str("timeframe multiplier must be greater than zero")
28 }
29}
30
31impl std::error::Error for TimeframeError {}
32
33impl Timeframe {
34 pub fn validate(self) -> Result<Self, TimeframeError> {
35 let value = match self {
36 Self::Second(v)
37 | Self::Minute(v)
38 | Self::Hour(v)
39 | Self::Day(v)
40 | Self::Week(v)
41 | Self::Month(v) => v,
42 };
43 (value > 0).then_some(self).ok_or(TimeframeError)
44 }
45
46 pub fn bucket_close(self, bucket_open_ts: i64) -> Option<i64> {
51 self.fixed_seconds()
52 .map(|secs| bucket_open_ts + secs as i64)
53 }
54
55 pub fn fixed_seconds(self) -> Option<u64> {
57 match self {
58 Self::Second(v) => Some(v as u64),
59 Self::Minute(v) => Some(v as u64 * 60),
60 Self::Hour(v) => Some(v as u64 * 3_600),
61 Self::Day(v) => Some(v as u64 * 86_400),
62 Self::Week(v) => Some(v as u64 * 604_800),
63 Self::Month(_) => None,
64 }
65 }
66
67 pub fn parse_str(value: &str) -> Option<Self> {
69 let trimmed = value.trim();
70 let unit = trimmed.chars().last()?;
71 let number = &trimmed[..trimmed.len().checked_sub(unit.len_utf8())?];
72 let multiplier = if number.is_empty() {
73 1
74 } else {
75 number.parse().ok()?
76 };
77 let timeframe = match unit {
78 's' | 'S' => Self::Second(multiplier),
79 'm' => Self::Minute(multiplier),
80 'h' | 'H' => Self::Hour(multiplier),
81 'd' | 'D' => Self::Day(multiplier),
82 'w' | 'W' => Self::Week(multiplier),
83 'M' => Self::Month(multiplier),
84 _ => return None,
85 };
86 timeframe.validate().ok()
87 }
88
89 pub(crate) fn bucket_start(self, timestamp: i64, utc_offset_seconds: i32) -> i64 {
90 let local = timestamp + i64::from(utc_offset_seconds);
91 let start_local = match self {
92 Self::Second(v) => fixed_bucket(local, i64::from(v)),
93 Self::Minute(v) => fixed_bucket(local, i64::from(v) * 60),
94 Self::Hour(v) => fixed_bucket(local, i64::from(v) * 3_600),
95 Self::Day(v) => fixed_bucket(local, i64::from(v) * SECONDS_PER_DAY),
96 Self::Week(v) => {
97 let width = i64::from(v) * 7;
98 let day = local.div_euclid(SECONDS_PER_DAY);
99 let start_day = (day + 3).div_euclid(width) * width - 3;
100 start_day * SECONDS_PER_DAY
101 }
102 Self::Month(v) => {
103 let day = local.div_euclid(SECONDS_PER_DAY);
104 let (year, month, _) = civil_from_days(day);
105 let month_index = i64::from(year) * 12 + i64::from(month) - 1;
106 let width = i64::from(v);
107 let start_index = month_index.div_euclid(width) * width;
108 let start_year = i32::try_from(start_index.div_euclid(12)).unwrap_or(1970);
109 let start_month = u32::try_from(start_index.rem_euclid(12)).unwrap_or(0) + 1;
110 days_from_civil(start_year, start_month, 1) * SECONDS_PER_DAY
111 }
112 };
113 start_local - i64::from(utc_offset_seconds)
114 }
115}
116
117fn fixed_bucket(timestamp: i64, width: i64) -> i64 {
118 timestamp.div_euclid(width) * width
119}
120
121fn civil_from_days(days: i64) -> (i32, u32, u32) {
123 let z = days + 719_468;
124 let era = z.div_euclid(146_097);
125 let doe = z - era * 146_097;
126 let yoe = (doe - doe / 1_460 + doe / 36_524 - doe / 146_096) / 365;
127 let mut year = yoe + era * 400;
128 let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
129 let mp = (5 * doy + 2) / 153;
130 let day = doy - (153 * mp + 2) / 5 + 1;
131 let month = mp + if mp < 10 { 3 } else { -9 };
132 year += i64::from(month <= 2);
133 (
134 i32::try_from(year).unwrap_or(1970),
135 u32::try_from(month).unwrap_or(1),
136 u32::try_from(day).unwrap_or(1),
137 )
138}
139
140fn days_from_civil(year: i32, month: u32, day: u32) -> i64 {
141 let adjusted_year = i64::from(year) - i64::from(month <= 2);
142 let era = adjusted_year.div_euclid(400);
143 let yoe = adjusted_year - era * 400;
144 let shifted_month = i64::from(month) + if month > 2 { -3 } else { 9 };
145 let doy = (153 * shifted_month + 2) / 5 + i64::from(day) - 1;
146 let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
147 era * 146_097 + doe - 719_468
148}
149
150impl fmt::Display for Timeframe {
151 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
152 match self {
153 Self::Second(v) => write!(f, "{v}s"),
154 Self::Minute(v) => write!(f, "{v}m"),
155 Self::Hour(v) => write!(f, "{v}h"),
156 Self::Day(v) => write!(f, "{v}d"),
157 Self::Week(v) => write!(f, "{v}w"),
158 Self::Month(v) => write!(f, "{v}M"),
159 }
160 }
161}
162
163#[derive(Debug, Clone, PartialEq)]
164pub struct ResamplerOutput {
165 pub completed_bar: Option<Bar>,
170 pub current_unconfirmed: Bar,
175 pub gap_buckets: u32,
182}
183
184#[derive(Debug, Clone)]
186pub struct BarResampler {
187 target_tf: Timeframe,
188 utc_offset_seconds: i32,
189 current_bucket: Option<Bar>,
190 bucket_start_ts: i64,
191}
192
193impl BarResampler {
194 pub fn new(target_tf: Timeframe) -> Result<Self, TimeframeError> {
195 Self::with_utc_offset(target_tf, 0)
196 }
197
198 pub fn with_utc_offset(
199 target_tf: Timeframe,
200 utc_offset_seconds: i32,
201 ) -> Result<Self, TimeframeError> {
202 Ok(Self {
203 target_tf: target_tf.validate()?,
204 utc_offset_seconds,
205 current_bucket: None,
206 bucket_start_ts: 0,
207 })
208 }
209
210 pub fn reset(&mut self) {
211 self.current_bucket = None;
212 self.bucket_start_ts = 0;
213 }
214
215 pub fn on_bar(&mut self, bar: &Bar) -> ResamplerOutput {
216 let bucket_start = self
217 .target_tf
218 .bucket_start(bar.timestamp, self.utc_offset_seconds);
219 let mut completed_bar = None;
220 let mut gap_buckets = 0u32;
221
222 if let Some(mut current) = self.current_bucket.take() {
223 if bucket_start != self.bucket_start_ts {
224 completed_bar = Some(current);
225 gap_buckets = self.gap_buckets_between(self.bucket_start_ts, bucket_start);
226 self.bucket_start_ts = bucket_start;
227 self.current_bucket = Some(start_bucket(bucket_start, bar));
228 } else {
229 current.high = current.high.max(bar.high);
230 current.low = current.low.min(bar.low);
231 current.close = bar.close;
232 current.volume += bar.volume;
233 self.current_bucket = Some(current);
234 }
235 } else {
236 self.bucket_start_ts = bucket_start;
237 self.current_bucket = Some(start_bucket(bucket_start, bar));
238 }
239
240 ResamplerOutput {
241 completed_bar,
242 current_unconfirmed: self.current_bucket.clone().expect("bucket was initialized"),
243 gap_buckets,
244 }
245 }
246
247 fn gap_buckets_between(&self, prev_start: i64, next_start: i64) -> u32 {
251 match self.target_tf.fixed_seconds() {
252 Some(width) if width > 0 => {
253 let delta_buckets = (next_start - prev_start) / (width as i64);
254 delta_buckets.saturating_sub(1).max(0) as u32
255 }
256 _ => 0,
257 }
258 }
259}
260
261#[derive(Debug, Clone)]
266pub struct ConfirmedResampler {
267 inner: BarResampler,
268}
269
270impl ConfirmedResampler {
271 pub fn new(target_tf: Timeframe) -> Result<Self, TimeframeError> {
272 Ok(Self {
273 inner: BarResampler::new(target_tf)?,
274 })
275 }
276
277 pub fn with_utc_offset(
278 target_tf: Timeframe,
279 utc_offset_seconds: i32,
280 ) -> Result<Self, TimeframeError> {
281 Ok(Self {
282 inner: BarResampler::with_utc_offset(target_tf, utc_offset_seconds)?,
283 })
284 }
285
286 pub fn reset(&mut self) {
287 self.inner.reset();
288 }
289
290 pub fn on_bar(&mut self, bar: &Bar) -> Option<Bar> {
293 self.inner.on_bar(bar).completed_bar
294 }
295}
296
297fn start_bucket(timestamp: i64, bar: &Bar) -> Bar {
298 Bar::new(
299 timestamp, bar.open, bar.high, bar.low, bar.close, bar.volume,
300 )
301}
302
303#[cfg(test)]
304mod tests {
305 use super::*;
306
307 #[test]
308 fn parsing_rejects_zero_and_supports_seconds() {
309 assert_eq!(Timeframe::parse_str("30s"), Some(Timeframe::Second(30)));
310 assert_eq!(Timeframe::parse_str("15m"), Some(Timeframe::Minute(15)));
311 assert_eq!(Timeframe::parse_str("1M"), Some(Timeframe::Month(1)));
312 assert_eq!(Timeframe::parse_str("0m"), None);
313 }
314
315 #[test]
316 fn resamples_fixed_intervals() {
317 let mut resampler = BarResampler::new(Timeframe::Minute(5)).unwrap();
318 for i in 0..5 {
319 let out = resampler.on_bar(&Bar::new(
320 i * 60,
321 100.0,
322 105.0,
323 95.0,
324 100.0 + i as f64,
325 100.0,
326 ));
327 assert!(out.completed_bar.is_none());
328 }
329 let out = resampler.on_bar(&Bar::new(300, 105.0, 110.0, 104.0, 108.0, 100.0));
330 let completed = out.completed_bar.unwrap();
331 assert_eq!(completed.timestamp, 0);
332 assert_eq!(completed.close, 104.0);
333 assert_eq!(completed.volume, 500.0);
334 }
335
336 #[test]
337 fn month_boundaries_are_calendar_aligned() {
338 let february = Timeframe::Month(1).bucket_start(1_706_745_600, 0);
339 let march = Timeframe::Month(1).bucket_start(1_709_251_200, 0);
340 assert_eq!(february, 1_706_745_600);
341 assert_eq!(march, 1_709_251_200);
342 assert_ne!(march - february, 30 * SECONDS_PER_DAY);
343 }
344
345 #[test]
346 fn negative_timestamps_use_euclidean_buckets() {
347 assert_eq!(Timeframe::Day(1).bucket_start(-1, 0), -86_400);
348 }
349
350 #[test]
351 fn bucket_close_matches_fixed_width_and_is_none_for_month() {
352 assert_eq!(Timeframe::Minute(5).bucket_close(0), Some(300));
353 assert_eq!(Timeframe::Day(1).bucket_close(0), Some(SECONDS_PER_DAY));
354 assert_eq!(Timeframe::Month(1).bucket_close(0), None);
355 }
356
357 #[test]
358 fn gap_buckets_is_zero_for_contiguous_bars() {
359 let mut resampler = BarResampler::new(Timeframe::Minute(5)).unwrap();
360 resampler.on_bar(&Bar::new(0, 100.0, 101.0, 99.0, 100.0, 10.0));
361 let out = resampler.on_bar(&Bar::new(300, 100.0, 101.0, 99.0, 100.0, 10.0));
362 assert_eq!(out.gap_buckets, 0);
363 }
364
365 #[test]
366 fn gap_buckets_reports_skipped_htf_buckets() {
367 let mut resampler = BarResampler::new(Timeframe::Minute(5)).unwrap();
368 resampler.on_bar(&Bar::new(0, 100.0, 101.0, 99.0, 100.0, 10.0));
369 let out = resampler.on_bar(&Bar::new(900, 100.0, 101.0, 99.0, 100.0, 10.0));
371 assert!(out.completed_bar.is_some());
372 assert_eq!(out.gap_buckets, 2);
373 }
374
375 #[test]
376 fn confirmed_resampler_never_exposes_unconfirmed_bucket() {
377 let mut confirmed = ConfirmedResampler::new(Timeframe::Minute(5)).unwrap();
378 for i in 0..5 {
379 let out = confirmed.on_bar(&Bar::new(
380 i * 60,
381 100.0,
382 105.0,
383 95.0,
384 100.0 + i as f64,
385 100.0,
386 ));
387 assert!(
388 out.is_none(),
389 "bucket must not repaint through ConfirmedResampler"
390 );
391 }
392 let closed = confirmed.on_bar(&Bar::new(300, 105.0, 110.0, 104.0, 108.0, 100.0));
393 let bar = closed.unwrap();
394 assert_eq!(bar.timestamp, 0);
395 assert_eq!(bar.close, 104.0);
396 }
397}