Skip to main content

kestrel_chartkit/
runner.rs

1//! Batch and replay execution over a full bar history.
2//!
3//! [`Indicator::on_bar`](crate::indicator::Indicator::on_bar) processes one bar and returns the
4//! latest value; consumers wanting a full, timestamp-aligned output series over a backfill
5//! (charting, backtesting, golden-fixture generation) previously had to hand-roll the loop. These
6//! helpers standardize it: deterministic (always start from [`Indicator::reset`]), timestamp-
7//! aligned (one entry per input bar, in order, `None` during warmup), and reproducible (pure
8//! function of the indicator + bar slice, safe to call repeatedly for replay).
9
10use crate::indicator::{Indicator, IndicatorOutput};
11use crate::model::{Bar, BarValidationError, SeriesCapabilities};
12
13/// One entry of a batch/replay output series: the source bar's timestamp paired with the
14/// indicator's output for that bar (`None` while still inside the warmup period).
15#[derive(Debug, Clone, PartialEq)]
16pub struct TimestampedOutput {
17    pub timestamp: i64,
18    pub output: Option<IndicatorOutput>,
19}
20
21/// Resets `indicator`, then feeds `bars` through it in order, collecting one [`TimestampedOutput`]
22/// per bar. Two calls with the same indicator type and `bars` slice always produce identical
23/// results (deterministic backfill / reproducible replay).
24pub fn run_batch<I: Indicator + ?Sized>(indicator: &mut I, bars: &[Bar]) -> Vec<TimestampedOutput> {
25    indicator.reset();
26    bars.iter()
27        .map(|bar| TimestampedOutput {
28            timestamp: bar.timestamp,
29            output: indicator.on_bar(bar),
30        })
31        .collect()
32}
33
34/// Like [`run_batch`], but validates each bar via
35/// [`Indicator::on_checked_bar`] and stops at the
36/// first invalid bar, returning the entries collected so far plus the validation error.
37pub fn run_batch_checked<I: Indicator + ?Sized>(
38    indicator: &mut I,
39    bars: &[Bar],
40) -> Result<Vec<TimestampedOutput>, (Vec<TimestampedOutput>, BarValidationError)> {
41    indicator.reset();
42    let mut series = Vec::with_capacity(bars.len());
43    for bar in bars {
44        match indicator.on_checked_bar(bar) {
45            Ok(output) => series.push(TimestampedOutput {
46                timestamp: bar.timestamp,
47                output,
48            }),
49            Err(err) => return Err((series, err)),
50        }
51    }
52    Ok(series)
53}
54
55/// Result of [`run_batch_with_applicability`]: the batch output series plus the applicability
56/// verdict for the indicator/series-capabilities pair it was computed for.
57#[derive(Debug, Clone, PartialEq)]
58pub struct BatchResult {
59    pub series: Vec<TimestampedOutput>,
60    pub applicability: crate::applicability::Applicability,
61}
62
63/// Like [`run_batch`], but also attaches an [`crate::applicability::Applicability`] verdict for
64/// the named indicator against `capabilities`.
65///
66/// `run_batch`/`run_batch_checked` are generic over `I: Indicator` and never see the indicator's
67/// registry name, so they cannot look up its [`crate::applicability::DataRequirements`]
68/// themselves — hence this separate function that takes `name` explicitly, rather than a change
69/// to either existing function's signature.
70///
71/// The series is always computed, even when the verdict is
72/// [`crate::applicability::Applicability::Unsuitable`] — callers may legitimately want to see a
73/// `Degraded` result, and a batch caller can still act on `Unsuitable` since it is always present
74/// on [`BatchResult`], not silently dropped. This is what avoids the "computes and hides the
75/// warning" failure mode the applicability check exists to prevent.
76/// Also tags every emitted [`IndicatorOutput`] with `capabilities` (see
77/// [`IndicatorOutput::series_capabilities`]) — this is the one place in the crate that already
78/// receives a `SeriesCapabilities` value alongside the indicator run, so it is where the
79/// origin gets attached rather than requiring every one of the ~90 `Indicator` implementations to
80/// do it themselves.
81pub fn run_batch_with_applicability<I: Indicator + ?Sized>(
82    name: &str,
83    indicator: &mut I,
84    bars: &[Bar],
85    capabilities: &SeriesCapabilities,
86) -> BatchResult {
87    let requirements = crate::applicability::data_requirements(name);
88    let applicability = crate::applicability::check_applicability(&requirements, capabilities);
89    let series = run_batch(indicator, bars)
90        .into_iter()
91        .map(|entry| TimestampedOutput {
92            timestamp: entry.timestamp,
93            output: entry
94                .output
95                .map(|output| output.with_capabilities(*capabilities)),
96        })
97        .collect();
98    BatchResult {
99        series,
100        applicability,
101    }
102}
103
104#[cfg(test)]
105mod tests {
106    use super::*;
107    use crate::indicator::moving_averages::SmaEngine;
108
109    /// Builds valid bars around each close (offset so `low = close + 100.0 - 1.0` stays positive).
110    fn sample_bars(closes: &[f64]) -> Vec<Bar> {
111        closes
112            .iter()
113            .enumerate()
114            .map(|(i, &c)| {
115                let c = c + 100.0;
116                Bar::new((i as i64) * 60, c, c + 1.0, c - 1.0, c, 100.0)
117            })
118            .collect()
119    }
120
121    #[test]
122    fn test_run_batch_is_timestamp_aligned_and_warmup_aware() {
123        let bars = sample_bars(&[1.0, 2.0, 3.0, 4.0, 5.0]);
124        let mut sma = SmaEngine::new(3);
125        let series = run_batch(&mut sma, &bars);
126
127        assert_eq!(series.len(), bars.len());
128        for (entry, bar) in series.iter().zip(&bars) {
129            assert_eq!(entry.timestamp, bar.timestamp);
130        }
131        // Warmup: SmaEngine needs 3 bars before it emits a value.
132        assert!(series[0].output.is_none());
133        assert!(series[1].output.is_none());
134        assert!(series[2].output.is_some());
135        assert_eq!(series[2].output.as_ref().unwrap().value, 102.0);
136        assert_eq!(series[4].output.as_ref().unwrap().value, 104.0);
137    }
138
139    #[test]
140    fn test_run_batch_is_deterministic_and_resets_prior_state() {
141        let bars = sample_bars(&[1.0, 2.0, 3.0, 4.0, 5.0]);
142        let mut sma = SmaEngine::new(3);
143
144        let first = run_batch(&mut sma, &bars);
145        // Re-running over the same indicator instance must reset first, so replay is idempotent.
146        let second = run_batch(&mut sma, &bars);
147        assert_eq!(first, second);
148    }
149
150    #[test]
151    fn test_run_batch_checked_stops_at_invalid_bar() {
152        let mut bars = sample_bars(&[10.0, 20.0]);
153        bars.push(Bar::new(120, f64::NAN, 1.0, -1.0, 1.0, 100.0));
154        bars.push(Bar::new(180, 3.0, 4.0, 2.0, 3.0, 100.0));
155
156        let mut sma = SmaEngine::new(2);
157        let result = run_batch_checked(&mut sma, &bars);
158        let (partial, err) = result.unwrap_err();
159        assert_eq!(partial.len(), 2);
160        assert_eq!(err, BarValidationError::NonFiniteValue);
161    }
162
163    fn real_volume_capabilities(volume: crate::model::VolumeKind) -> SeriesCapabilities {
164        use crate::model::{
165            ContinuityKind, LiquidityTier, PriceAdjustment, Provenance, SessionKind,
166        };
167        SeriesCapabilities {
168            volume,
169            trade_direction: false,
170            session: SessionKind::Regular,
171            continuity: ContinuityKind::SingleContract,
172            price_adjustment: PriceAdjustment::Raw,
173            provenance: Provenance::Exchange,
174            liquidity_tier: LiquidityTier::Deep,
175        }
176    }
177
178    #[test]
179    fn test_run_batch_with_applicability_is_applicable_with_real_volume() {
180        let bars = sample_bars(&[1.0, 2.0, 3.0, 4.0, 5.0]);
181        let mut sma = SmaEngine::new(3);
182        let capabilities = real_volume_capabilities(crate::model::VolumeKind::RealTurnover);
183
184        let result = run_batch_with_applicability("vwap", &mut sma, &bars, &capabilities);
185
186        assert_eq!(result.series.len(), bars.len());
187        assert_eq!(
188            result.applicability,
189            crate::applicability::Applicability::Applicable
190        );
191    }
192
193    #[test]
194    fn test_run_batch_with_applicability_flags_unsuitable_but_still_computes() {
195        let bars = sample_bars(&[1.0, 2.0, 3.0, 4.0, 5.0]);
196        let mut sma = SmaEngine::new(3);
197        let capabilities = real_volume_capabilities(crate::model::VolumeKind::Tick);
198
199        let result = run_batch_with_applicability("vwap", &mut sma, &bars, &capabilities);
200
201        // The series is still computed even though the verdict is Unsuitable.
202        assert_eq!(result.series.len(), bars.len());
203        assert!(matches!(
204            result.applicability,
205            crate::applicability::Applicability::Unsuitable { .. }
206        ));
207    }
208
209    #[test]
210    fn test_run_batch_with_applicability_tags_every_output_with_capabilities() {
211        let bars = sample_bars(&[1.0, 2.0, 3.0, 4.0, 5.0]);
212        let mut sma = SmaEngine::new(3);
213        let capabilities = real_volume_capabilities(crate::model::VolumeKind::RealTurnover);
214
215        let result = run_batch_with_applicability("vwap", &mut sma, &bars, &capabilities);
216
217        // Every emitted output (i.e. every entry past warmup) carries the capabilities it was
218        // computed against — this is what makes pivots_structure/zigzag/zigzag_advanced/
219        // pivot_sets (and every other generic `Indicator`) traceable to their source series
220        // without each of them needing its own `series_capabilities` plumbing.
221        for entry in &result.series {
222            if let Some(output) = &entry.output {
223                assert_eq!(output.series_capabilities, Some(capabilities));
224            }
225        }
226        assert!(result.series.iter().any(|e| e.output.is_some()));
227    }
228
229    #[test]
230    fn test_run_batch_leaves_capabilities_unset() {
231        let bars = sample_bars(&[1.0, 2.0, 3.0, 4.0, 5.0]);
232        let mut sma = SmaEngine::new(3);
233
234        let series = run_batch(&mut sma, &bars);
235
236        // Plain `run_batch` never sees a `SeriesCapabilities` value, so it cannot attach one —
237        // callers without that information get `None`, same as a direct `on_bar` call.
238        for entry in &series {
239            if let Some(output) = &entry.output {
240                assert_eq!(output.series_capabilities, None);
241            }
242        }
243    }
244}