1use std::collections::{BTreeMap, BTreeSet, VecDeque};
6
7use chrono::NaiveDateTime;
8
9use crate::data_feed::{MarketEvent, TimestampBatch};
10
11use super::{PriceBasis, SeriesId, SeriesRequirement, StrategyRequirements, WarmupRequirement};
12
13pub const MAX_RETAINED_BARS: usize = 1_000_000;
14
15#[derive(Debug, Clone, Copy, PartialEq, Eq)]
17pub enum MissingIntervalPolicy {
18 Skip,
19 Reject,
20}
21
22#[derive(Debug, Clone, PartialEq, Eq)]
24pub struct BarSeriesSpec {
25 requirement: SeriesRequirement,
26 retained_bars: usize,
27 alignment_offset_seconds: i64,
28 missing_interval: MissingIntervalPolicy,
29}
30
31impl BarSeriesSpec {
32 pub fn new(
33 requirement: SeriesRequirement,
34 retained_bars: usize,
35 alignment_offset_seconds: i32,
36 missing_interval: MissingIntervalPolicy,
37 ) -> Result<Self, SeriesError> {
38 if retained_bars == 0 {
39 return Err(SeriesError::ZeroRetention {
40 series_id: requirement.id().clone(),
41 });
42 }
43 if retained_bars > MAX_RETAINED_BARS {
44 return Err(SeriesError::RetentionTooLarge {
45 series_id: requirement.id().clone(),
46 retained_bars,
47 maximum: MAX_RETAINED_BARS,
48 });
49 }
50 let required_bars = requirement.warmup().required_bars();
51 if retained_bars < required_bars {
52 return Err(SeriesError::RetentionBelowWarmup {
53 series_id: requirement.id().clone(),
54 retained_bars,
55 required_bars,
56 });
57 }
58 let duration = i64::try_from(requirement.timeframe().duration_seconds())
59 .expect("fixed timeframe duration always fits i64");
60 let alignment_offset_seconds = i64::from(alignment_offset_seconds).rem_euclid(duration);
61 Ok(Self {
62 requirement,
63 retained_bars,
64 alignment_offset_seconds,
65 missing_interval,
66 })
67 }
68
69 pub fn requirement(&self) -> &SeriesRequirement {
70 &self.requirement
71 }
72
73 pub fn retained_bars(&self) -> usize {
74 self.retained_bars
75 }
76
77 pub fn alignment_offset_seconds(&self) -> i64 {
78 self.alignment_offset_seconds
79 }
80
81 pub fn missing_interval(&self) -> MissingIntervalPolicy {
82 self.missing_interval
83 }
84}
85
86#[derive(Debug, Clone, PartialEq)]
88pub struct ClosedBar {
89 series_id: SeriesId,
90 symbol: String,
91 open_time: NaiveDateTime,
92 close_time: NaiveDateTime,
93 open: f64,
94 high: f64,
95 low: f64,
96 close: f64,
97 tick_count: Option<u64>,
98}
99
100impl ClosedBar {
101 #[cfg(test)]
102 pub(crate) fn for_test(
103 series_id: SeriesId,
104 symbol: impl Into<String>,
105 open_time: NaiveDateTime,
106 close_time: NaiveDateTime,
107 high: f64,
108 low: f64,
109 ) -> Self {
110 Self {
111 series_id,
112 symbol: symbol.into(),
113 open_time,
114 close_time,
115 open: low,
116 high,
117 low,
118 close: high,
119 tick_count: Some(1),
120 }
121 }
122
123 pub fn series_id(&self) -> &SeriesId {
124 &self.series_id
125 }
126
127 pub fn symbol(&self) -> &str {
128 &self.symbol
129 }
130
131 pub fn open_time(&self) -> NaiveDateTime {
132 self.open_time
133 }
134
135 pub fn close_time(&self) -> NaiveDateTime {
136 self.close_time
137 }
138
139 pub fn open(&self) -> f64 {
140 self.open
141 }
142
143 pub fn high(&self) -> f64 {
144 self.high
145 }
146
147 pub fn low(&self) -> f64 {
148 self.low
149 }
150
151 pub fn close(&self) -> f64 {
152 self.close
153 }
154
155 pub fn tick_count(&self) -> Option<u64> {
156 self.tick_count
157 }
158}
159
160#[derive(Debug, Clone, Copy)]
162pub struct BarWindow<'a> {
163 older: &'a [ClosedBar],
164 newer: &'a [ClosedBar],
165}
166
167impl<'a> BarWindow<'a> {
168 pub fn len(&self) -> usize {
169 self.older.len() + self.newer.len()
170 }
171
172 pub fn is_empty(&self) -> bool {
173 self.older.is_empty() && self.newer.is_empty()
174 }
175
176 pub fn latest(&self) -> Option<&'a ClosedBar> {
177 self.newer.last().or_else(|| self.older.last())
178 }
179
180 pub fn iter(&self) -> impl Iterator<Item = &'a ClosedBar> {
181 self.older.iter().chain(self.newer.iter())
182 }
183}
184
185#[derive(Debug, Clone, Copy, PartialEq, Eq)]
187pub struct SeriesWarmupState {
188 required: WarmupRequirement,
189 available_bars: usize,
190}
191
192impl SeriesWarmupState {
193 pub fn required(self) -> WarmupRequirement {
194 self.required
195 }
196
197 pub fn available_bars(self) -> usize {
198 self.available_bars
199 }
200
201 pub fn is_ready(self) -> bool {
202 self.available_bars >= self.required.required_bars()
203 }
204}
205
206pub trait HistoricalSeriesView {
208 fn latest_bar(&self, id: &SeriesId) -> Result<Option<&ClosedBar>, SeriesViewError>;
209 fn bars(&self, id: &SeriesId, count: usize) -> Result<BarWindow<'_>, SeriesViewError>;
210 fn warmup(&self, id: &SeriesId) -> Result<SeriesWarmupState, SeriesViewError>;
211}
212
213#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
215pub enum SeriesError {
216 #[error("series '{series_id}' retained-bar capacity must be greater than zero")]
217 ZeroRetention { series_id: SeriesId },
218 #[error("series '{series_id}' retained-bar capacity {retained_bars} exceeds maximum {maximum}")]
219 RetentionTooLarge {
220 series_id: SeriesId,
221 retained_bars: usize,
222 maximum: usize,
223 },
224 #[error(
225 "series '{series_id}' retained-bar capacity {retained_bars} is below warmup {required_bars}"
226 )]
227 RetentionBelowWarmup {
228 series_id: SeriesId,
229 retained_bars: usize,
230 required_bars: usize,
231 },
232 #[error("series ID '{series_id}' is configured more than once")]
233 DuplicateSeriesId { series_id: SeriesId },
234 #[error("batch at {batch_ts} contains an event at {event_ts}")]
235 BatchTimestampMismatch {
236 batch_ts: NaiveDateTime,
237 event_ts: NaiveDateTime,
238 },
239 #[error("duplicate event ordering metadata ({series_rank}, {row_sequence}) at {timestamp}")]
240 DuplicateOrderingMetadata {
241 timestamp: NaiveDateTime,
242 series_rank: u32,
243 row_sequence: u64,
244 },
245 #[error("primary tick for '{symbol}' moved backwards from {previous} to {current}")]
246 TimestampRegression {
247 symbol: String,
248 previous: NaiveDateTime,
249 current: NaiveDateTime,
250 },
251 #[error("series '{series_id}' cannot represent a bucket boundary for {timestamp}")]
252 BoundaryOverflow {
253 series_id: SeriesId,
254 timestamp: NaiveDateTime,
255 },
256 #[error("series '{series_id}' has one or more empty intervals before {next_open}")]
257 MissingInterval {
258 series_id: SeriesId,
259 previous_close: NaiveDateTime,
260 next_open: NaiveDateTime,
261 },
262 #[error("series '{series_id}' tick count overflowed")]
263 TickCountOverflow { series_id: SeriesId },
264 #[error("series '{series_id}' completed-bar count overflowed")]
265 CompletedBarCountOverflow { series_id: SeriesId },
266 #[error("series '{series_id}' received both ticks and stored bars")]
267 MixedSeriesInput { series_id: SeriesId },
268 #[error("stored bar for '{symbol}' at {timestamp} does not declare its timeframe")]
269 StoredBarWithoutTimeframe {
270 symbol: String,
271 timestamp: NaiveDateTime,
272 },
273 #[error(
274 "stored {timeframe_seconds}s bar for '{symbol}' at {timestamp} matches no declared series"
275 )]
276 StoredBarUnmatched {
277 symbol: String,
278 timeframe_seconds: u64,
279 timestamp: NaiveDateTime,
280 },
281 #[error(
282 "stored bar for series '{series_id}' at {timestamp} does not start on a bucket boundary; the bucket opens at {bucket_open}"
283 )]
284 StoredBarMisaligned {
285 series_id: SeriesId,
286 timestamp: NaiveDateTime,
287 bucket_open: NaiveDateTime,
288 },
289 #[error("stored bar for series '{series_id}' at {timestamp} is invalid: {reason}")]
290 InvalidStoredBar {
291 series_id: SeriesId,
292 timestamp: NaiveDateTime,
293 reason: &'static str,
294 },
295 #[error(
296 "series '{series_id}' received a second stored bar for the bucket opening at {timestamp}"
297 )]
298 DuplicateStoredBar {
299 series_id: SeriesId,
300 timestamp: NaiveDateTime,
301 },
302}
303
304#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
306pub enum SeriesViewError {
307 #[error("unknown historical series '{series_id}'")]
308 UnknownSeries { series_id: SeriesId },
309}
310
311#[derive(Debug, Clone)]
312struct OpenBar {
313 open_time: NaiveDateTime,
314 close_time: NaiveDateTime,
315 open: f64,
316 high: f64,
317 low: f64,
318 close: f64,
319 tick_count: Option<u64>,
320}
321
322impl OpenBar {
323 fn new(open_time: NaiveDateTime, close_time: NaiveDateTime, price: f64) -> Self {
324 Self {
325 open_time,
326 close_time,
327 open: price,
328 high: price,
329 low: price,
330 close: price,
331 tick_count: Some(1),
332 }
333 }
334
335 fn update(&mut self, series_id: &SeriesId, price: f64) -> Result<(), SeriesError> {
336 self.high = self.high.max(price);
337 self.low = self.low.min(price);
338 self.close = price;
339 self.tick_count = Some(
340 self.tick_count
341 .expect("tick-built bars always have a count")
342 .checked_add(1)
343 .ok_or_else(|| SeriesError::TickCountOverflow {
344 series_id: series_id.clone(),
345 })?,
346 );
347 Ok(())
348 }
349}
350
351#[derive(Debug, Clone, Copy, PartialEq, Eq)]
353enum SeriesInput {
354 Ticks,
355 StoredBars,
356}
357
358#[derive(Debug, Clone)]
359struct SeriesState {
360 spec: BarSeriesSpec,
361 open: Option<OpenBar>,
362 closed: VecDeque<ClosedBar>,
363 completed_bars: usize,
364 input: Option<SeriesInput>,
365}
366
367impl SeriesState {
368 fn new(spec: BarSeriesSpec) -> Self {
369 Self {
370 closed: VecDeque::with_capacity(spec.retained_bars),
371 spec,
372 open: None,
373 completed_bars: 0,
374 input: None,
375 }
376 }
377
378 fn duration_seconds(&self) -> i64 {
379 i64::try_from(self.spec.requirement.timeframe().duration_seconds())
380 .expect("fixed timeframe duration always fits i64")
381 }
382}
383
384#[derive(Debug)]
385struct BatchSeriesState {
386 open: Option<OpenBar>,
387 completed_bars: usize,
388 emitted: Vec<ClosedBar>,
389 input: Option<SeriesInput>,
390}
391
392impl BatchSeriesState {
393 fn from_committed(state: &SeriesState) -> Self {
394 Self {
395 open: state.open.clone(),
396 completed_bars: state.completed_bars,
397 emitted: Vec::new(),
398 input: state.input,
399 }
400 }
401
402 fn accept_input(
403 &mut self,
404 spec: &BarSeriesSpec,
405 input: SeriesInput,
406 ) -> Result<(), SeriesError> {
407 match self.input {
408 Some(existing) if existing != input => Err(SeriesError::MixedSeriesInput {
409 series_id: spec.requirement.id().clone(),
410 }),
411 _ => {
412 self.input = Some(input);
413 Ok(())
414 }
415 }
416 }
417
418 fn apply_tick(
419 &mut self,
420 spec: &BarSeriesSpec,
421 timestamp: NaiveDateTime,
422 bid: f64,
423 ask: f64,
424 ) -> Result<(), SeriesError> {
425 let price = match spec.requirement.price_basis() {
426 PriceBasis::Bid => bid,
427 PriceBasis::Ask => ask,
428 PriceBasis::Mid => bid + (ask - bid) / 2.0,
429 };
430 let duration = i64::try_from(spec.requirement.timeframe().duration_seconds())
431 .expect("fixed timeframe duration always fits i64");
432 let (open_time, close_time) = bucket_bounds(
433 spec.requirement.id(),
434 timestamp,
435 duration,
436 spec.alignment_offset_seconds,
437 )?;
438
439 let Some(current) = self.open.as_mut() else {
440 self.open = Some(OpenBar::new(open_time, close_time, price));
441 return Ok(());
442 };
443 if open_time == current.open_time {
444 current.update(spec.requirement.id(), price)?;
445 return Ok(());
446 }
447
448 self.close_open_bar(spec, open_time)?;
449 self.open = Some(OpenBar::new(open_time, close_time, price));
450 Ok(())
451 }
452
453 fn apply_stored_bar(
455 &mut self,
456 spec: &BarSeriesSpec,
457 bar: StoredBar,
458 ) -> Result<(), SeriesError> {
459 let series_id = spec.requirement.id();
460 let duration = i64::try_from(spec.requirement.timeframe().duration_seconds())
461 .expect("fixed timeframe duration always fits i64");
462 let (open_time, close_time) =
463 bucket_bounds(series_id, bar.ts, duration, spec.alignment_offset_seconds)?;
464 if open_time != bar.ts {
465 return Err(SeriesError::StoredBarMisaligned {
466 series_id: series_id.clone(),
467 timestamp: bar.ts,
468 bucket_open: open_time,
469 });
470 }
471 let invalid = |reason| SeriesError::InvalidStoredBar {
472 series_id: series_id.clone(),
473 timestamp: bar.ts,
474 reason,
475 };
476 if ![bar.open, bar.high, bar.low, bar.close]
477 .iter()
478 .all(|value| value.is_finite())
479 {
480 return Err(invalid("prices must be finite"));
481 }
482 if bar.low > bar.open.min(bar.close) || bar.high < bar.open.max(bar.close) {
483 return Err(invalid("open and close must lie within low and high"));
484 }
485 if bar.tick_count == Some(0) {
486 return Err(invalid("zero is not a valid tick count"));
487 }
488 let tick_count = bar.tick_count;
489 if let Some(current) = self.open.as_ref() {
490 if current.open_time == open_time {
491 return Err(SeriesError::DuplicateStoredBar {
492 series_id: series_id.clone(),
493 timestamp: bar.ts,
494 });
495 }
496 self.close_open_bar(spec, open_time)?;
497 }
498 self.open = Some(OpenBar {
499 open_time,
500 close_time,
501 open: bar.open,
502 high: bar.high,
503 low: bar.low,
504 close: bar.close,
505 tick_count,
506 });
507 Ok(())
508 }
509
510 fn close_open_bar(
512 &mut self,
513 spec: &BarSeriesSpec,
514 next_open: NaiveDateTime,
515 ) -> Result<(), SeriesError> {
516 let current = self
517 .open
518 .as_ref()
519 .expect("a held bar exists when a later bucket arrives");
520 if spec.missing_interval == MissingIntervalPolicy::Reject && next_open > current.close_time
521 {
522 return Err(SeriesError::MissingInterval {
523 series_id: spec.requirement.id().clone(),
524 previous_close: current.close_time,
525 next_open,
526 });
527 }
528 let completed = self.open.take().expect("held bar checked above");
529 self.completed_bars = self.completed_bars.checked_add(1).ok_or_else(|| {
530 SeriesError::CompletedBarCountOverflow {
531 series_id: spec.requirement.id().clone(),
532 }
533 })?;
534 self.emitted.push(ClosedBar {
535 series_id: spec.requirement.id().clone(),
536 symbol: spec.requirement.symbol().to_string(),
537 open_time: completed.open_time,
538 close_time: completed.close_time,
539 open: completed.open,
540 high: completed.high,
541 low: completed.low,
542 close: completed.close,
543 tick_count: completed.tick_count,
544 });
545 Ok(())
546 }
547}
548
549#[derive(Debug, Clone, Copy)]
551struct StoredBar {
552 ts: NaiveDateTime,
553 open: f64,
554 high: f64,
555 low: f64,
556 close: f64,
557 tick_count: Option<u64>,
558}
559
560#[derive(Debug)]
562pub struct MultiTimeframeSeries {
563 series: BTreeMap<SeriesId, SeriesState>,
564 last_source_ts: BTreeMap<String, NaiveDateTime>,
565}
566
567impl MultiTimeframeSeries {
568 pub fn new(specs: Vec<BarSeriesSpec>) -> Result<Self, SeriesError> {
569 let mut series = BTreeMap::new();
570 for spec in specs {
571 let id = spec.requirement.id().clone();
572 if series.insert(id.clone(), SeriesState::new(spec)).is_some() {
573 return Err(SeriesError::DuplicateSeriesId { series_id: id });
574 }
575 }
576 Ok(Self {
577 series,
578 last_source_ts: BTreeMap::new(),
579 })
580 }
581
582 pub fn on_batch(&mut self, batch: &TimestampBatch) -> Result<Vec<ClosedBar>, SeriesError> {
583 validate_batch(batch)?;
584 let mut staged_series = self
585 .series
586 .iter()
587 .map(|(id, state)| (id.clone(), BatchSeriesState::from_committed(state)))
588 .collect::<BTreeMap<_, _>>();
589 let mut staged_source_ts = self.last_source_ts.clone();
590 self.preflight_batch(batch, &mut staged_series, &mut staged_source_ts)?;
591
592 let mut emitted = Vec::new();
593 for (id, mut staged) in staged_series {
594 let state = self
595 .series
596 .get_mut(&id)
597 .expect("staged series originates from committed state");
598 state.open = staged.open;
599 state.completed_bars = staged.completed_bars;
600 state.input = staged.input;
601 for closed in staged.emitted.drain(..) {
602 if state.closed.len() == state.spec.retained_bars {
603 state.closed.pop_front();
604 }
605 state.closed.push_back(closed.clone());
606 emitted.push(closed);
607 }
608 }
609 self.last_source_ts = staged_source_ts;
610 emitted.sort_by(|left, right| {
611 left.close_time
612 .cmp(&right.close_time)
613 .then_with(|| {
614 self.series[left.series_id()]
615 .duration_seconds()
616 .cmp(&self.series[right.series_id()].duration_seconds())
617 })
618 .then_with(|| left.series_id.cmp(&right.series_id))
619 });
620 Ok(emitted)
621 }
622
623 pub fn latest_bar(&self, series: &SeriesId) -> Result<Option<&ClosedBar>, SeriesViewError> {
624 HistoricalSeriesView::latest_bar(self, series)
625 }
626
627 pub fn bars(&self, series: &SeriesId, count: usize) -> Result<BarWindow<'_>, SeriesViewError> {
628 HistoricalSeriesView::bars(self, series, count)
629 }
630
631 pub fn warmup(&self, series: &SeriesId) -> Result<SeriesWarmupState, SeriesViewError> {
632 HistoricalSeriesView::warmup(self, series)
633 }
634
635 pub fn warmup_complete(
636 &self,
637 requirements: &StrategyRequirements,
638 ) -> Result<bool, SeriesViewError> {
639 for requirement in requirements.series() {
640 self.state(requirement.id())?;
641 }
642 Ok(requirements
643 .warmup_complete(|id| self.series.get(id).map_or(0, |state| state.completed_bars)))
644 }
645
646 fn preflight_batch(
647 &self,
648 batch: &TimestampBatch,
649 staged_series: &mut BTreeMap<SeriesId, BatchSeriesState>,
650 staged_source_ts: &mut BTreeMap<String, NaiveDateTime>,
651 ) -> Result<(), SeriesError> {
652 let mut ordered = batch.events.iter().collect::<Vec<_>>();
653 ordered.sort_by_key(|event| (event.metadata.series_rank, event.metadata.row_sequence));
654
655 for feed_event in ordered {
656 if !feed_event.metadata.roles.primary {
657 continue;
658 }
659 let symbol = feed_event.event.symbol();
660 let ts = feed_event.event.ts();
661 if let Some(previous) = staged_source_ts.get(symbol)
662 && ts < *previous
663 {
664 return Err(SeriesError::TimestampRegression {
665 symbol: symbol.to_owned(),
666 previous: *previous,
667 current: ts,
668 });
669 }
670 staged_source_ts.insert(symbol.to_owned(), ts);
671
672 match &feed_event.event {
673 MarketEvent::Tick { bid, ask, .. } => {
674 if feed_event.event.to_valid_quote().is_none() {
675 continue;
676 }
677 for (id, state) in self
678 .series
679 .iter()
680 .filter(|(_, state)| state.spec.requirement.symbol() == symbol)
681 {
682 let staged = staged_series
683 .get_mut(id)
684 .expect("staged series originates from committed state");
685 staged.accept_input(&state.spec, SeriesInput::Ticks)?;
686 staged.apply_tick(&state.spec, ts, *bid, *ask)?;
687 }
688 }
689 MarketEvent::Bar {
690 open,
691 high,
692 low,
693 close,
694 timeframe_seconds,
695 tick_count,
696 ..
697 } => {
698 if !self
699 .series
700 .values()
701 .any(|state| state.spec.requirement.symbol() == symbol)
702 {
703 continue;
704 }
705 let timeframe_seconds = timeframe_seconds.ok_or_else(|| {
706 SeriesError::StoredBarWithoutTimeframe {
707 symbol: symbol.to_owned(),
708 timestamp: ts,
709 }
710 })?;
711 let bar = StoredBar {
712 ts,
713 open: *open,
714 high: *high,
715 low: *low,
716 close: *close,
717 tick_count: *tick_count,
718 };
719 let mut matched = false;
720 for (id, state) in self.series.iter().filter(|(_, state)| {
721 state.spec.requirement.symbol() == symbol
722 && state.spec.requirement.timeframe().duration_seconds()
723 == timeframe_seconds
724 }) {
725 matched = true;
726 let staged = staged_series
727 .get_mut(id)
728 .expect("staged series originates from committed state");
729 staged.accept_input(&state.spec, SeriesInput::StoredBars)?;
730 staged.apply_stored_bar(&state.spec, bar)?;
731 }
732 if !matched {
733 return Err(SeriesError::StoredBarUnmatched {
734 symbol: symbol.to_owned(),
735 timeframe_seconds,
736 timestamp: ts,
737 });
738 }
739 }
740 }
741 }
742 Ok(())
743 }
744
745 fn state(&self, id: &SeriesId) -> Result<&SeriesState, SeriesViewError> {
746 self.series
747 .get(id)
748 .ok_or_else(|| SeriesViewError::UnknownSeries {
749 series_id: id.clone(),
750 })
751 }
752}
753
754impl HistoricalSeriesView for MultiTimeframeSeries {
755 fn latest_bar(&self, id: &SeriesId) -> Result<Option<&ClosedBar>, SeriesViewError> {
756 Ok(self.state(id)?.closed.back())
757 }
758
759 fn bars(&self, id: &SeriesId, count: usize) -> Result<BarWindow<'_>, SeriesViewError> {
760 let state = self.state(id)?;
761 let (older, newer) = state.closed.as_slices();
762 let available = count.min(state.closed.len());
763 let skip = state.closed.len() - available;
764 if skip < older.len() {
765 Ok(BarWindow {
766 older: &older[skip..],
767 newer,
768 })
769 } else {
770 Ok(BarWindow {
771 older: &older[older.len()..],
772 newer: &newer[skip - older.len()..],
773 })
774 }
775 }
776
777 fn warmup(&self, id: &SeriesId) -> Result<SeriesWarmupState, SeriesViewError> {
778 let state = self.state(id)?;
779 Ok(SeriesWarmupState {
780 required: state.spec.requirement.warmup(),
781 available_bars: state.completed_bars,
782 })
783 }
784}
785
786fn validate_batch(batch: &TimestampBatch) -> Result<(), SeriesError> {
787 let mut ordering = BTreeSet::new();
788 for feed_event in &batch.events {
789 let event_ts = feed_event.event.ts();
790 if event_ts != batch.ts {
791 return Err(SeriesError::BatchTimestampMismatch {
792 batch_ts: batch.ts,
793 event_ts,
794 });
795 }
796 let key = (
797 feed_event.metadata.series_rank,
798 feed_event.metadata.row_sequence,
799 );
800 if !ordering.insert(key) {
801 return Err(SeriesError::DuplicateOrderingMetadata {
802 timestamp: batch.ts,
803 series_rank: key.0,
804 row_sequence: key.1,
805 });
806 }
807 }
808 Ok(())
809}
810
811fn bucket_bounds(
815 series_id: &SeriesId,
816 timestamp: NaiveDateTime,
817 duration: i64,
818 offset: i64,
819) -> Result<(NaiveDateTime, NaiveDateTime), SeriesError> {
820 let spec = data_preprocess::resample::BucketSpec::new(duration, offset).ok_or_else(|| {
821 SeriesError::BoundaryOverflow {
822 series_id: series_id.clone(),
823 timestamp,
824 }
825 })?;
826 data_preprocess::resample::bucket_bounds(timestamp, spec).ok_or_else(|| {
827 SeriesError::BoundaryOverflow {
828 series_id: series_id.clone(),
829 timestamp,
830 }
831 })
832}