Skip to main content

qs_backtest/strategy/
series.rs

1//! Causal fixed-duration closed bars derived from historical tick batches.
2
3use std::collections::{BTreeMap, BTreeSet, VecDeque};
4
5use chrono::{DateTime, NaiveDateTime};
6
7use crate::data_feed::{MarketEvent, TimestampBatch};
8
9use super::{PriceBasis, SeriesId, SeriesRequirement, StrategyRequirements, WarmupRequirement};
10
11pub const MAX_RETAINED_BARS: usize = 1_000_000;
12
13/// Missing fixed-duration buckets between observed ticks.
14#[derive(Debug, Clone, Copy, PartialEq, Eq)]
15pub enum MissingIntervalPolicy {
16    Skip,
17    Reject,
18}
19
20/// Operational configuration for one tick-derived closed-bar series.
21#[derive(Debug, Clone, PartialEq, Eq)]
22pub struct BarSeriesSpec {
23    requirement: SeriesRequirement,
24    retained_bars: usize,
25    alignment_offset_seconds: i64,
26    missing_interval: MissingIntervalPolicy,
27}
28
29impl BarSeriesSpec {
30    pub fn new(
31        requirement: SeriesRequirement,
32        retained_bars: usize,
33        alignment_offset_seconds: i32,
34        missing_interval: MissingIntervalPolicy,
35    ) -> Result<Self, SeriesError> {
36        if retained_bars == 0 {
37            return Err(SeriesError::ZeroRetention {
38                series_id: requirement.id().clone(),
39            });
40        }
41        if retained_bars > MAX_RETAINED_BARS {
42            return Err(SeriesError::RetentionTooLarge {
43                series_id: requirement.id().clone(),
44                retained_bars,
45                maximum: MAX_RETAINED_BARS,
46            });
47        }
48        let required_bars = requirement.warmup().required_bars();
49        if retained_bars < required_bars {
50            return Err(SeriesError::RetentionBelowWarmup {
51                series_id: requirement.id().clone(),
52                retained_bars,
53                required_bars,
54            });
55        }
56        let duration = i64::try_from(requirement.timeframe().duration_seconds())
57            .expect("fixed timeframe duration always fits i64");
58        let alignment_offset_seconds = i64::from(alignment_offset_seconds).rem_euclid(duration);
59        Ok(Self {
60            requirement,
61            retained_bars,
62            alignment_offset_seconds,
63            missing_interval,
64        })
65    }
66
67    pub fn requirement(&self) -> &SeriesRequirement {
68        &self.requirement
69    }
70
71    pub fn retained_bars(&self) -> usize {
72        self.retained_bars
73    }
74
75    pub fn alignment_offset_seconds(&self) -> i64 {
76        self.alignment_offset_seconds
77    }
78
79    pub fn missing_interval(&self) -> MissingIntervalPolicy {
80        self.missing_interval
81    }
82}
83
84/// One immutable nonempty bar visible at its exclusive close boundary.
85#[derive(Debug, Clone, PartialEq)]
86pub struct ClosedBar {
87    series_id: SeriesId,
88    symbol: String,
89    open_time: NaiveDateTime,
90    close_time: NaiveDateTime,
91    open: f64,
92    high: f64,
93    low: f64,
94    close: f64,
95    tick_count: u64,
96}
97
98impl ClosedBar {
99    pub fn series_id(&self) -> &SeriesId {
100        &self.series_id
101    }
102
103    pub fn symbol(&self) -> &str {
104        &self.symbol
105    }
106
107    pub fn open_time(&self) -> NaiveDateTime {
108        self.open_time
109    }
110
111    pub fn close_time(&self) -> NaiveDateTime {
112        self.close_time
113    }
114
115    pub fn open(&self) -> f64 {
116        self.open
117    }
118
119    pub fn high(&self) -> f64 {
120        self.high
121    }
122
123    pub fn low(&self) -> f64 {
124        self.low
125    }
126
127    pub fn close(&self) -> f64 {
128        self.close
129    }
130
131    pub fn tick_count(&self) -> u64 {
132        self.tick_count
133    }
134}
135
136/// Allocation-free suffix view over a possibly wrapped retained deque.
137#[derive(Debug, Clone, Copy)]
138pub struct BarWindow<'a> {
139    older: &'a [ClosedBar],
140    newer: &'a [ClosedBar],
141}
142
143impl<'a> BarWindow<'a> {
144    pub fn len(&self) -> usize {
145        self.older.len() + self.newer.len()
146    }
147
148    pub fn is_empty(&self) -> bool {
149        self.older.is_empty() && self.newer.is_empty()
150    }
151
152    pub fn latest(&self) -> Option<&'a ClosedBar> {
153        self.newer.last().or_else(|| self.older.last())
154    }
155
156    pub fn iter(&self) -> impl Iterator<Item = &'a ClosedBar> {
157        self.older.iter().chain(self.newer.iter())
158    }
159}
160
161/// Readiness facts for one configured series.
162#[derive(Debug, Clone, Copy, PartialEq, Eq)]
163pub struct SeriesWarmupState {
164    required: WarmupRequirement,
165    available_bars: usize,
166}
167
168impl SeriesWarmupState {
169    pub fn required(self) -> WarmupRequirement {
170        self.required
171    }
172
173    pub fn available_bars(self) -> usize {
174        self.available_bars
175    }
176
177    pub fn is_ready(self) -> bool {
178        self.available_bars >= self.required.required_bars()
179    }
180}
181
182/// Read-only causal history exposed to analyzers and strategies.
183pub trait HistoricalSeriesView {
184    fn latest_bar(&self, id: &SeriesId) -> Result<Option<&ClosedBar>, SeriesViewError>;
185    fn bars(&self, id: &SeriesId, count: usize) -> Result<BarWindow<'_>, SeriesViewError>;
186    fn warmup(&self, id: &SeriesId) -> Result<SeriesWarmupState, SeriesViewError>;
187}
188
189/// Errors returned while constructing or updating historical series.
190#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
191pub enum SeriesError {
192    #[error("series '{series_id}' retained-bar capacity must be greater than zero")]
193    ZeroRetention { series_id: SeriesId },
194    #[error("series '{series_id}' retained-bar capacity {retained_bars} exceeds maximum {maximum}")]
195    RetentionTooLarge {
196        series_id: SeriesId,
197        retained_bars: usize,
198        maximum: usize,
199    },
200    #[error(
201        "series '{series_id}' retained-bar capacity {retained_bars} is below warmup {required_bars}"
202    )]
203    RetentionBelowWarmup {
204        series_id: SeriesId,
205        retained_bars: usize,
206        required_bars: usize,
207    },
208    #[error("series ID '{series_id}' is configured more than once")]
209    DuplicateSeriesId { series_id: SeriesId },
210    #[error("batch at {batch_ts} contains an event at {event_ts}")]
211    BatchTimestampMismatch {
212        batch_ts: NaiveDateTime,
213        event_ts: NaiveDateTime,
214    },
215    #[error("duplicate event ordering metadata ({series_rank}, {row_sequence}) at {timestamp}")]
216    DuplicateOrderingMetadata {
217        timestamp: NaiveDateTime,
218        series_rank: u32,
219        row_sequence: u64,
220    },
221    #[error("primary tick for '{symbol}' moved backwards from {previous} to {current}")]
222    TimestampRegression {
223        symbol: String,
224        previous: NaiveDateTime,
225        current: NaiveDateTime,
226    },
227    #[error("series '{series_id}' cannot represent a bucket boundary for {timestamp}")]
228    BoundaryOverflow {
229        series_id: SeriesId,
230        timestamp: NaiveDateTime,
231    },
232    #[error("series '{series_id}' has one or more empty intervals before {next_open}")]
233    MissingInterval {
234        series_id: SeriesId,
235        previous_close: NaiveDateTime,
236        next_open: NaiveDateTime,
237    },
238    #[error("series '{series_id}' tick count overflowed")]
239    TickCountOverflow { series_id: SeriesId },
240    #[error("series '{series_id}' completed-bar count overflowed")]
241    CompletedBarCountOverflow { series_id: SeriesId },
242}
243
244/// Errors returned by read-only series lookup.
245#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
246pub enum SeriesViewError {
247    #[error("unknown historical series '{series_id}'")]
248    UnknownSeries { series_id: SeriesId },
249}
250
251#[derive(Debug, Clone)]
252struct OpenBar {
253    open_time: NaiveDateTime,
254    close_time: NaiveDateTime,
255    open: f64,
256    high: f64,
257    low: f64,
258    close: f64,
259    tick_count: u64,
260}
261
262impl OpenBar {
263    fn new(open_time: NaiveDateTime, close_time: NaiveDateTime, price: f64) -> Self {
264        Self {
265            open_time,
266            close_time,
267            open: price,
268            high: price,
269            low: price,
270            close: price,
271            tick_count: 1,
272        }
273    }
274
275    fn update(&mut self, series_id: &SeriesId, price: f64) -> Result<(), SeriesError> {
276        self.high = self.high.max(price);
277        self.low = self.low.min(price);
278        self.close = price;
279        self.tick_count =
280            self.tick_count
281                .checked_add(1)
282                .ok_or_else(|| SeriesError::TickCountOverflow {
283                    series_id: series_id.clone(),
284                })?;
285        Ok(())
286    }
287}
288
289#[derive(Debug, Clone)]
290struct SeriesState {
291    spec: BarSeriesSpec,
292    open: Option<OpenBar>,
293    closed: VecDeque<ClosedBar>,
294    completed_bars: usize,
295}
296
297impl SeriesState {
298    fn new(spec: BarSeriesSpec) -> Self {
299        Self {
300            closed: VecDeque::with_capacity(spec.retained_bars),
301            spec,
302            open: None,
303            completed_bars: 0,
304        }
305    }
306
307    fn duration_seconds(&self) -> i64 {
308        i64::try_from(self.spec.requirement.timeframe().duration_seconds())
309            .expect("fixed timeframe duration always fits i64")
310    }
311}
312
313#[derive(Debug)]
314struct BatchSeriesState {
315    open: Option<OpenBar>,
316    completed_bars: usize,
317    emitted: Vec<ClosedBar>,
318}
319
320impl BatchSeriesState {
321    fn from_committed(state: &SeriesState) -> Self {
322        Self {
323            open: state.open.clone(),
324            completed_bars: state.completed_bars,
325            emitted: Vec::new(),
326        }
327    }
328
329    fn apply_tick(
330        &mut self,
331        spec: &BarSeriesSpec,
332        timestamp: NaiveDateTime,
333        bid: f64,
334        ask: f64,
335    ) -> Result<(), SeriesError> {
336        let price = match spec.requirement.price_basis() {
337            PriceBasis::Bid => bid,
338            PriceBasis::Ask => ask,
339            PriceBasis::Mid => bid + (ask - bid) / 2.0,
340        };
341        let duration = i64::try_from(spec.requirement.timeframe().duration_seconds())
342            .expect("fixed timeframe duration always fits i64");
343        let (open_time, close_time) = bucket_bounds(
344            spec.requirement.id(),
345            timestamp,
346            duration,
347            spec.alignment_offset_seconds,
348        )?;
349
350        let Some(current) = self.open.as_mut() else {
351            self.open = Some(OpenBar::new(open_time, close_time, price));
352            return Ok(());
353        };
354        if open_time == current.open_time {
355            current.update(spec.requirement.id(), price)?;
356            return Ok(());
357        }
358
359        if spec.missing_interval == MissingIntervalPolicy::Reject && open_time > current.close_time
360        {
361            return Err(SeriesError::MissingInterval {
362                series_id: spec.requirement.id().clone(),
363                previous_close: current.close_time,
364                next_open: open_time,
365            });
366        }
367
368        let completed = self
369            .open
370            .take()
371            .expect("open bar exists after transition validation");
372        self.completed_bars = self.completed_bars.checked_add(1).ok_or_else(|| {
373            SeriesError::CompletedBarCountOverflow {
374                series_id: spec.requirement.id().clone(),
375            }
376        })?;
377        self.emitted.push(ClosedBar {
378            series_id: spec.requirement.id().clone(),
379            symbol: spec.requirement.symbol().to_string(),
380            open_time: completed.open_time,
381            close_time: completed.close_time,
382            open: completed.open,
383            high: completed.high,
384            low: completed.low,
385            close: completed.close,
386            tick_count: completed.tick_count,
387        });
388        self.open = Some(OpenBar::new(open_time, close_time, price));
389        Ok(())
390    }
391}
392
393/// Bounded causal closed-bar state for several symbols and timeframes.
394#[derive(Debug)]
395pub struct MultiTimeframeSeries {
396    series: BTreeMap<SeriesId, SeriesState>,
397    last_source_ts: BTreeMap<String, NaiveDateTime>,
398}
399
400impl MultiTimeframeSeries {
401    pub fn new(specs: Vec<BarSeriesSpec>) -> Result<Self, SeriesError> {
402        let mut series = BTreeMap::new();
403        for spec in specs {
404            let id = spec.requirement.id().clone();
405            if series.insert(id.clone(), SeriesState::new(spec)).is_some() {
406                return Err(SeriesError::DuplicateSeriesId { series_id: id });
407            }
408        }
409        Ok(Self {
410            series,
411            last_source_ts: BTreeMap::new(),
412        })
413    }
414
415    pub fn on_batch(&mut self, batch: &TimestampBatch) -> Result<Vec<ClosedBar>, SeriesError> {
416        validate_batch(batch)?;
417        let mut staged_series = self
418            .series
419            .iter()
420            .map(|(id, state)| (id.clone(), BatchSeriesState::from_committed(state)))
421            .collect::<BTreeMap<_, _>>();
422        let mut staged_source_ts = self.last_source_ts.clone();
423        self.preflight_batch(batch, &mut staged_series, &mut staged_source_ts)?;
424
425        let mut emitted = Vec::new();
426        for (id, mut staged) in staged_series {
427            let state = self
428                .series
429                .get_mut(&id)
430                .expect("staged series originates from committed state");
431            state.open = staged.open;
432            state.completed_bars = staged.completed_bars;
433            for closed in staged.emitted.drain(..) {
434                if state.closed.len() == state.spec.retained_bars {
435                    state.closed.pop_front();
436                }
437                state.closed.push_back(closed.clone());
438                emitted.push(closed);
439            }
440        }
441        self.last_source_ts = staged_source_ts;
442        emitted.sort_by(|left, right| {
443            left.close_time
444                .cmp(&right.close_time)
445                .then_with(|| {
446                    self.series[left.series_id()]
447                        .duration_seconds()
448                        .cmp(&self.series[right.series_id()].duration_seconds())
449                })
450                .then_with(|| left.series_id.cmp(&right.series_id))
451        });
452        Ok(emitted)
453    }
454
455    pub fn latest_bar(&self, series: &SeriesId) -> Result<Option<&ClosedBar>, SeriesViewError> {
456        HistoricalSeriesView::latest_bar(self, series)
457    }
458
459    pub fn bars(&self, series: &SeriesId, count: usize) -> Result<BarWindow<'_>, SeriesViewError> {
460        HistoricalSeriesView::bars(self, series, count)
461    }
462
463    pub fn warmup(&self, series: &SeriesId) -> Result<SeriesWarmupState, SeriesViewError> {
464        HistoricalSeriesView::warmup(self, series)
465    }
466
467    pub fn warmup_complete(
468        &self,
469        requirements: &StrategyRequirements,
470    ) -> Result<bool, SeriesViewError> {
471        for requirement in requirements.series() {
472            self.state(requirement.id())?;
473        }
474        Ok(requirements
475            .warmup_complete(|id| self.series.get(id).map_or(0, |state| state.completed_bars)))
476    }
477
478    fn preflight_batch(
479        &self,
480        batch: &TimestampBatch,
481        staged_series: &mut BTreeMap<SeriesId, BatchSeriesState>,
482        staged_source_ts: &mut BTreeMap<String, NaiveDateTime>,
483    ) -> Result<(), SeriesError> {
484        let mut ordered = batch.events.iter().collect::<Vec<_>>();
485        ordered.sort_by_key(|event| (event.metadata.series_rank, event.metadata.row_sequence));
486
487        for feed_event in ordered {
488            if !feed_event.metadata.roles.primary {
489                continue;
490            }
491            let MarketEvent::Tick {
492                symbol,
493                ts,
494                bid,
495                ask,
496            } = &feed_event.event
497            else {
498                continue;
499            };
500            if let Some(previous) = staged_source_ts.get(symbol)
501                && *ts < *previous
502            {
503                return Err(SeriesError::TimestampRegression {
504                    symbol: symbol.clone(),
505                    previous: *previous,
506                    current: *ts,
507                });
508            }
509            staged_source_ts.insert(symbol.clone(), *ts);
510
511            if feed_event.event.to_valid_quote().is_none() {
512                continue;
513            }
514            for (id, state) in self
515                .series
516                .iter()
517                .filter(|(_, state)| state.spec.requirement.symbol() == symbol)
518            {
519                staged_series
520                    .get_mut(id)
521                    .expect("staged series originates from committed state")
522                    .apply_tick(&state.spec, *ts, *bid, *ask)?;
523            }
524        }
525        Ok(())
526    }
527
528    fn state(&self, id: &SeriesId) -> Result<&SeriesState, SeriesViewError> {
529        self.series
530            .get(id)
531            .ok_or_else(|| SeriesViewError::UnknownSeries {
532                series_id: id.clone(),
533            })
534    }
535}
536
537impl HistoricalSeriesView for MultiTimeframeSeries {
538    fn latest_bar(&self, id: &SeriesId) -> Result<Option<&ClosedBar>, SeriesViewError> {
539        Ok(self.state(id)?.closed.back())
540    }
541
542    fn bars(&self, id: &SeriesId, count: usize) -> Result<BarWindow<'_>, SeriesViewError> {
543        let state = self.state(id)?;
544        let (older, newer) = state.closed.as_slices();
545        let available = count.min(state.closed.len());
546        let skip = state.closed.len() - available;
547        if skip < older.len() {
548            Ok(BarWindow {
549                older: &older[skip..],
550                newer,
551            })
552        } else {
553            Ok(BarWindow {
554                older: &older[older.len()..],
555                newer: &newer[skip - older.len()..],
556            })
557        }
558    }
559
560    fn warmup(&self, id: &SeriesId) -> Result<SeriesWarmupState, SeriesViewError> {
561        let state = self.state(id)?;
562        Ok(SeriesWarmupState {
563            required: state.spec.requirement.warmup(),
564            available_bars: state.completed_bars,
565        })
566    }
567}
568
569fn validate_batch(batch: &TimestampBatch) -> Result<(), SeriesError> {
570    let mut ordering = BTreeSet::new();
571    for feed_event in &batch.events {
572        let event_ts = feed_event.event.ts();
573        if event_ts != batch.ts {
574            return Err(SeriesError::BatchTimestampMismatch {
575                batch_ts: batch.ts,
576                event_ts,
577            });
578        }
579        let key = (
580            feed_event.metadata.series_rank,
581            feed_event.metadata.row_sequence,
582        );
583        if !ordering.insert(key) {
584            return Err(SeriesError::DuplicateOrderingMetadata {
585                timestamp: batch.ts,
586                series_rank: key.0,
587                row_sequence: key.1,
588            });
589        }
590    }
591    Ok(())
592}
593
594fn bucket_bounds(
595    series_id: &SeriesId,
596    timestamp: NaiveDateTime,
597    duration: i64,
598    offset: i64,
599) -> Result<(NaiveDateTime, NaiveDateTime), SeriesError> {
600    let timestamp_seconds = i128::from(timestamp.and_utc().timestamp());
601    let duration = i128::from(duration);
602    let offset = i128::from(offset);
603    let bucket_index = (timestamp_seconds - offset).div_euclid(duration);
604    let open_seconds = bucket_index
605        .checked_mul(duration)
606        .and_then(|value| value.checked_add(offset))
607        .ok_or_else(|| SeriesError::BoundaryOverflow {
608            series_id: series_id.clone(),
609            timestamp,
610        })?;
611    let close_seconds =
612        open_seconds
613            .checked_add(duration)
614            .ok_or_else(|| SeriesError::BoundaryOverflow {
615                series_id: series_id.clone(),
616                timestamp,
617            })?;
618    let open_seconds = i64::try_from(open_seconds).map_err(|_| SeriesError::BoundaryOverflow {
619        series_id: series_id.clone(),
620        timestamp,
621    })?;
622    let close_seconds =
623        i64::try_from(close_seconds).map_err(|_| SeriesError::BoundaryOverflow {
624            series_id: series_id.clone(),
625            timestamp,
626        })?;
627    let open_time = DateTime::from_timestamp(open_seconds, 0)
628        .map(|value| value.naive_utc())
629        .ok_or_else(|| SeriesError::BoundaryOverflow {
630            series_id: series_id.clone(),
631            timestamp,
632        })?;
633    let close_time = DateTime::from_timestamp(close_seconds, 0)
634        .map(|value| value.naive_utc())
635        .ok_or_else(|| SeriesError::BoundaryOverflow {
636            series_id: series_id.clone(),
637            timestamp,
638        })?;
639    Ok((open_time, close_time))
640}