Skip to main content

qs_backtest/strategy/
series.rs

1//! Causal fixed-duration closed bars built from historical tick batches or accepted from stored bars.
2//!
3//! A series is fed by one kind of primary input. Ticks accumulate into the bucket that contains them. A stored bar is stamped at its bucket's open time, as the resampler writes it, and is held as that bucket's complete contents; like a tick-built bucket, it becomes visible only when a later bucket for the same symbol arrives, so a strategy never sees a bar before its close.
4
5use 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/// Missing fixed-duration buckets between observed ticks.
16#[derive(Debug, Clone, Copy, PartialEq, Eq)]
17pub enum MissingIntervalPolicy {
18    Skip,
19    Reject,
20}
21
22/// Operational configuration for one tick-derived closed-bar series.
23#[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/// One immutable nonempty bar visible at its exclusive close boundary.
87#[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/// Allocation-free suffix view over a possibly wrapped retained deque.
161#[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/// Readiness facts for one configured series.
186#[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
206/// Read-only causal history exposed to analyzers and strategies.
207pub 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/// Errors returned while constructing or updating historical series.
214#[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/// Errors returned by read-only series lookup.
305#[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/// The one kind of primary input a series has accepted.
352#[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    /// Hold one stored bar as the complete contents of its bucket, completing the previously held bar first.
454    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    /// Complete the held bar because a later bucket opening at `next_open` has arrived.
511    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/// Fields of one primary stored-bar event, as the series reads them.
550#[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/// Bounded causal closed-bar state for several symbols and timeframes.
561#[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
811/// Half-open `[open, close)` bucket containing `timestamp`.
812///
813/// The arithmetic is shared with the stored-bar resampler so that a bar built during replay and the same bar written to storage are identical.
814fn 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}