Skip to main content

kestrel_chartkit/
lifecycle.rs

1//! Bar lifecycle events and rollback-safe, idempotent recomputation.
2//!
3//! [`Indicator::on_bar`](crate::indicator::Indicator::on_bar) takes a single incoming bar and
4//! permanently advances the indicator's state — it has no notion of a bar that is still forming
5//! versus one that has closed, and no way to reprocess a bar without double-counting it.
6//!
7//! [`BarLifecycle`] names the four events a feed can emit for a given timestamp. [`LifecycleRunner`]
8//! wraps any `Indicator + Clone` and makes [`BarLifecycle::Update`]/[`BarLifecycle::Correction`]
9//! safe to call repeatedly: it always recomputes from the last *confirmed* checkpoint rather than
10//! mutating forward, so re-delivering the same revised bar is idempotent and never double-applies
11//! a still-forming bar's data.
12//!
13//! [`LifecycleRunner`] keeps two checkpoints for the last confirmed bar: the state immediately
14//! *before* it was applied and the state immediately *after*. A [`BarLifecycle::Correction`] (or a
15//! repeated [`BarLifecycle::Confirmed`] for the same timestamp) restores the pre-bar checkpoint and
16//! applies the revised bar exactly once, so the corrected bar *replaces* the original rather than
17//! stacking on top of it. Only the most recently confirmed bar can be corrected this way — a
18//! correction or an out-of-order confirmation targeting an older timestamp is rejected with
19//! [`LifecycleError`] rather than silently mutating an unrelated state.
20//!
21//! Requiring `Clone` is a deliberate, honest scope limit: it is the only mechanism that works
22//! generically across arbitrary indicator internals without adding a second, hand-written
23//! snapshot/restore method to every `Indicator` impl. Most indicator structs in this crate do not
24//! yet derive `Clone`; adding it is a mechanical, low-risk follow-up left to indicator authors as
25//! they adopt `LifecycleRunner`, not a blanket change bundled into this module.
26
27use crate::checkpoint::CheckpointStore;
28use crate::indicator::{Indicator, IndicatorOutput};
29use crate::model::Bar;
30use std::fmt;
31
32/// A bar-lifecycle event for a single timestamp.
33#[derive(Debug, Clone, PartialEq)]
34pub enum BarLifecycle {
35    /// The first tick of a new bar/timestamp.
36    Open(Bar),
37    /// A subsequent, still-unconfirmed tick for the same timestamp as the last `Open`/`Update`:
38    /// the bar is still forming and its values will be revised.
39    Update(Bar),
40    /// The bar for this timestamp is final and will not be revised again.
41    Confirmed(Bar),
42    /// A feed correction to the most recently confirmed bar (e.g. an exchange restatement).
43    /// Only single-level undo is supported: correcting a bar older than the last confirmed one
44    /// is not representable by [`LifecycleRunner`], which keeps just one pre/post checkpoint pair
45    /// and rejects such a correction with [`LifecycleError::CorrectionTimestampMismatch`].
46    Correction(Bar),
47}
48
49/// An error produced by [`LifecycleRunner::on_event`] when an event cannot be applied without
50/// losing or misrepresenting state.
51#[derive(Debug, Clone, Copy, PartialEq, Eq)]
52pub enum LifecycleError {
53    /// A [`BarLifecycle::Correction`] targeted a timestamp other than the last confirmed bar's.
54    /// `LifecycleRunner` only supports single-level undo: it can replace the most recently
55    /// confirmed bar, not an older one.
56    CorrectionTimestampMismatch {
57        /// The timestamp of the last confirmed bar, if any.
58        last_confirmed_timestamp: Option<i64>,
59        /// The timestamp the correction targeted.
60        attempted_timestamp: i64,
61    },
62    /// A [`BarLifecycle::Confirmed`] event arrived for a timestamp strictly older than the last
63    /// confirmed bar. `LifecycleRunner` cannot represent this: only the most recently confirmed
64    /// bar can be replaced (via a matching `Confirmed` or `Correction`).
65    NonMonotonicConfirmation {
66        /// The timestamp of the last confirmed bar.
67        last_confirmed_timestamp: i64,
68        /// The timestamp the out-of-order confirmation targeted.
69        attempted_timestamp: i64,
70    },
71}
72
73impl fmt::Display for LifecycleError {
74    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
75        match self {
76            LifecycleError::CorrectionTimestampMismatch {
77                last_confirmed_timestamp,
78                attempted_timestamp,
79            } => write!(
80                f,
81                "correction for timestamp {attempted_timestamp} does not match the last \
82                 confirmed timestamp {last_confirmed_timestamp:?}; only the most recently \
83                 confirmed bar can be corrected"
84            ),
85            LifecycleError::NonMonotonicConfirmation {
86                last_confirmed_timestamp,
87                attempted_timestamp,
88            } => write!(
89                f,
90                "confirmed timestamp {attempted_timestamp} is older than the last confirmed \
91                 timestamp {last_confirmed_timestamp}; only the most recently confirmed bar can \
92                 be replaced"
93            ),
94        }
95    }
96}
97
98impl std::error::Error for LifecycleError {}
99
100impl BarLifecycle {
101    pub fn bar(&self) -> &Bar {
102        match self {
103            BarLifecycle::Open(b)
104            | BarLifecycle::Update(b)
105            | BarLifecycle::Confirmed(b)
106            | BarLifecycle::Correction(b) => b,
107        }
108    }
109}
110
111/// Wraps an `Indicator + Clone` so [`BarLifecycle::Open`]/[`BarLifecycle::Update`] events can be
112/// replayed any number of times without permanently mutating state, and only
113/// [`BarLifecycle::Confirmed`]/[`BarLifecycle::Correction`] commit a new checkpoint.
114pub struct LifecycleRunner<I: Indicator + Clone> {
115    live: I,
116    /// State immediately before the last confirmed bar was applied (the rewind point for a
117    /// correction/replacement of that bar).
118    pre_confirmed: CheckpointStore<I>,
119    /// State immediately after the last confirmed bar (the rewind point for `Open`/`Update`).
120    post_confirmed: CheckpointStore<I>,
121}
122
123impl<I: Indicator + Clone> LifecycleRunner<I> {
124    pub fn new(indicator: I) -> Self {
125        Self {
126            live: indicator,
127            pre_confirmed: CheckpointStore::new(),
128            post_confirmed: CheckpointStore::new(),
129        }
130    }
131
132    /// The indicator state as of the last confirmed/corrected bar (never a tentative
133    /// open/update). `None` before the first confirmation.
134    pub fn confirmed_indicator(&self) -> Option<&I> {
135        self.post_confirmed.latest().map(|c| &c.state)
136    }
137
138    /// Processes one lifecycle event and returns the resulting output.
139    ///
140    /// `Open`/`Update` always start from the last confirmed checkpoint (or a freshly reset
141    /// indicator if none exists yet) before applying `bar`, so repeated calls for the same
142    /// still-forming timestamp are idempotent and never stack on top of a prior tentative
143    /// application.
144    ///
145    /// `Confirmed` for a new (strictly later, or first-ever) timestamp advances both checkpoints.
146    /// `Confirmed` or `Correction` for the same timestamp as the last confirmed bar *replaces* it:
147    /// the pre-bar checkpoint is restored and the revised bar is applied exactly once, instead of
148    /// stacking on top of the original bar's effect. A `Correction`/out-of-order `Confirmed` for
149    /// any other timestamp is rejected with [`LifecycleError`], since only single-level undo is
150    /// representable.
151    pub fn on_event(
152        &mut self,
153        event: BarLifecycle,
154    ) -> Result<Option<IndicatorOutput>, LifecycleError> {
155        match event {
156            BarLifecycle::Open(bar) | BarLifecycle::Update(bar) => {
157                self.rewind_to_last_confirmed();
158                Ok(self.live.on_bar(&bar))
159            }
160            BarLifecycle::Confirmed(bar) => match self.last_confirmed_timestamp() {
161                Some(ts) if bar.timestamp == ts => Ok(self.replace_last_confirmed(&bar)),
162                Some(ts) if bar.timestamp < ts => Err(LifecycleError::NonMonotonicConfirmation {
163                    last_confirmed_timestamp: ts,
164                    attempted_timestamp: bar.timestamp,
165                }),
166                _ => {
167                    self.rewind_to_last_confirmed();
168                    let pre_state = self.live.clone();
169                    let output = self.live.on_bar(&bar);
170                    self.pre_confirmed.save(&pre_state, bar.timestamp);
171                    self.post_confirmed.save(&self.live, bar.timestamp);
172                    Ok(output)
173                }
174            },
175            BarLifecycle::Correction(bar) => match self.last_confirmed_timestamp() {
176                Some(ts) if bar.timestamp == ts => Ok(self.replace_last_confirmed(&bar)),
177                last_ts => Err(LifecycleError::CorrectionTimestampMismatch {
178                    last_confirmed_timestamp: last_ts,
179                    attempted_timestamp: bar.timestamp,
180                }),
181            },
182        }
183    }
184
185    /// Restores the checkpoint from immediately before the last confirmed bar and applies `bar`
186    /// (which shares that bar's timestamp) exactly once, replacing it.
187    fn replace_last_confirmed(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
188        if !self.pre_confirmed.restore_into(&mut self.live) {
189            self.live.reset();
190        }
191        let output = self.live.on_bar(bar);
192        self.post_confirmed.save(&self.live, bar.timestamp);
193        output
194    }
195
196    fn last_confirmed_timestamp(&self) -> Option<i64> {
197        self.post_confirmed.latest().map(|c| c.timestamp)
198    }
199
200    fn rewind_to_last_confirmed(&mut self) {
201        if !self.post_confirmed.restore_into(&mut self.live) {
202            self.live.reset();
203        }
204    }
205}
206
207#[cfg(test)]
208mod tests {
209    use super::*;
210    use crate::indicator::moving_averages::EmaEngine;
211
212    // SmaEngine does not derive Clone (see module docs), so tests use a minimal hand-rolled Clone
213    // indicator instead of pulling SmaEngine into the Clone surface. Being a cumulative sum, it
214    // makes double-application (the bug in finding 01) visible: a `LastCloseEngine`-style
215    // overwrite would hide it.
216    #[derive(Clone)]
217    struct SumEngine {
218        sum: f64,
219    }
220    impl Indicator for SumEngine {
221        fn name(&self) -> &str {
222            "sum"
223        }
224        fn reset(&mut self) {
225            self.sum = 0.0;
226        }
227        fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
228            self.sum += bar.close;
229            Some(IndicatorOutput::new(self.sum))
230        }
231    }
232
233    #[test]
234    fn test_repeated_update_is_idempotent() {
235        let mut runner = LifecycleRunner::new(SumEngine { sum: 0.0 });
236        let bar = Bar::new(1_000, 100.0, 101.0, 99.0, 100.0, 10.0);
237
238        let first_update = runner.on_event(BarLifecycle::Open(bar.clone())).unwrap();
239        let second_update = runner.on_event(BarLifecycle::Update(bar.clone())).unwrap();
240        // Re-delivering the same still-forming bar must not accumulate.
241        assert_eq!(first_update.unwrap().value, second_update.unwrap().value);
242
243        let confirmed = runner
244            .on_event(BarLifecycle::Confirmed(bar.clone()))
245            .unwrap();
246        assert_eq!(confirmed.unwrap().value, 100.0);
247        assert_eq!(runner.confirmed_indicator().unwrap().sum, 100.0);
248
249        // A later bar accumulates on top of the confirmed checkpoint, not the discarded update.
250        let next_bar = Bar::new(1_060, 100.0, 102.0, 100.0, 101.0, 10.0);
251        let next = runner.on_event(BarLifecycle::Confirmed(next_bar)).unwrap();
252        assert_eq!(next.unwrap().value, 201.0);
253    }
254
255    /// Minimal reproduction from finding 01: a correction must *replace* the last confirmed bar's
256    /// contribution, not add to it. With the pre-fix implementation this asserted `203.5`.
257    #[test]
258    fn test_correction_replaces_rather_than_accumulates() {
259        let mut runner = LifecycleRunner::new(SumEngine { sum: 0.0 });
260        let bar = Bar::new(1_000, 100.0, 101.0, 99.0, 100.0, 10.0);
261        runner.on_event(BarLifecycle::Confirmed(bar)).unwrap();
262
263        let corrected_bar = Bar::new(1_000, 100.0, 101.0, 99.0, 103.5, 10.0);
264        let corrected = runner
265            .on_event(BarLifecycle::Correction(corrected_bar))
266            .unwrap();
267        assert_eq!(corrected.unwrap().value, 103.5);
268        assert_eq!(runner.confirmed_indicator().unwrap().sum, 103.5);
269    }
270
271    /// A correction to the second of two confirmed bars must replace only that bar; the first
272    /// bar's contribution must survive untouched.
273    #[test]
274    fn test_correction_of_second_bar_preserves_first_bar() {
275        let mut runner = LifecycleRunner::new(SumEngine { sum: 0.0 });
276        runner
277            .on_event(BarLifecycle::Confirmed(Bar::new(
278                1_000, 100.0, 101.0, 99.0, 100.0, 10.0,
279            )))
280            .unwrap();
281        runner
282            .on_event(BarLifecycle::Confirmed(Bar::new(
283                1_060, 100.0, 102.0, 100.0, 101.0, 10.0,
284            )))
285            .unwrap();
286
287        let corrected = runner
288            .on_event(BarLifecycle::Correction(Bar::new(
289                1_060, 100.0, 102.0, 100.0, 102.0, 10.0,
290            )))
291            .unwrap();
292        assert_eq!(corrected.unwrap().value, 202.0);
293        assert_eq!(runner.confirmed_indicator().unwrap().sum, 202.0);
294    }
295
296    /// A correction whose timestamp does not match the last confirmed bar cannot be represented
297    /// (only single-level undo is supported) and must be rejected explicitly rather than silently
298    /// mutating an unrelated checkpoint.
299    #[test]
300    fn test_correction_with_mismatched_timestamp_is_rejected() {
301        let mut runner = LifecycleRunner::new(SumEngine { sum: 0.0 });
302        runner
303            .on_event(BarLifecycle::Confirmed(Bar::new(
304                1_000, 100.0, 101.0, 99.0, 100.0, 10.0,
305            )))
306            .unwrap();
307        runner
308            .on_event(BarLifecycle::Confirmed(Bar::new(
309                1_060, 100.0, 102.0, 100.0, 101.0, 10.0,
310            )))
311            .unwrap();
312
313        // Targets the first bar's timestamp, not the last confirmed one.
314        let err = runner
315            .on_event(BarLifecycle::Correction(Bar::new(
316                1_000, 100.0, 101.0, 99.0, 999.0, 10.0,
317            )))
318            .unwrap_err();
319        assert_eq!(
320            err,
321            LifecycleError::CorrectionTimestampMismatch {
322                last_confirmed_timestamp: Some(1_060),
323                attempted_timestamp: 1_000,
324            }
325        );
326        // State is unchanged after a rejected correction.
327        assert_eq!(runner.confirmed_indicator().unwrap().sum, 201.0);
328    }
329
330    /// A `Confirmed` event repeated for the same timestamp as the last confirmed bar behaves like
331    /// a correction: it replaces rather than double-applies.
332    #[test]
333    fn test_repeated_confirmed_for_same_timestamp_replaces() {
334        let mut runner = LifecycleRunner::new(SumEngine { sum: 0.0 });
335        let bar = Bar::new(1_000, 100.0, 101.0, 99.0, 100.0, 10.0);
336        runner.on_event(BarLifecycle::Confirmed(bar)).unwrap();
337
338        let re_confirmed = Bar::new(1_000, 100.0, 101.0, 99.0, 105.0, 10.0);
339        let output = runner
340            .on_event(BarLifecycle::Confirmed(re_confirmed))
341            .unwrap();
342        assert_eq!(output.unwrap().value, 105.0);
343        assert_eq!(runner.confirmed_indicator().unwrap().sum, 105.0);
344    }
345
346    /// A `Confirmed` event for a timestamp older than the last confirmed bar is out-of-order and
347    /// not representable; it must be rejected instead of silently corrupting state.
348    #[test]
349    fn test_confirmed_with_older_timestamp_is_rejected() {
350        let mut runner = LifecycleRunner::new(SumEngine { sum: 0.0 });
351        runner
352            .on_event(BarLifecycle::Confirmed(Bar::new(
353                1_060, 100.0, 102.0, 100.0, 101.0, 10.0,
354            )))
355            .unwrap();
356
357        let err = runner
358            .on_event(BarLifecycle::Confirmed(Bar::new(
359                1_000, 100.0, 101.0, 99.0, 100.0, 10.0,
360            )))
361            .unwrap_err();
362        assert_eq!(
363            err,
364            LifecycleError::NonMonotonicConfirmation {
365                last_confirmed_timestamp: 1_060,
366                attempted_timestamp: 1_000,
367            }
368        );
369        assert_eq!(runner.confirmed_indicator().unwrap().sum, 101.0);
370    }
371
372    /// Same correction-replaces-rather-than-accumulates scenario, but against a real recursive
373    /// indicator (`EmaEngine`) rather than the artificial `SumEngine`, so the fix is confirmed
374    /// against production indicator state and not just a test double.
375    #[test]
376    fn test_correction_replaces_for_real_indicator() {
377        let mut runner = LifecycleRunner::new(EmaEngine::new(2));
378        let k = 2.0 / 3.0;
379
380        runner
381            .on_event(BarLifecycle::Confirmed(Bar::new(
382                1_000, 100.0, 101.0, 99.0, 100.0, 10.0,
383            )))
384            .unwrap();
385        let confirmed = runner
386            .on_event(BarLifecycle::Confirmed(Bar::new(
387                1_060, 100.0, 106.0, 100.0, 106.0, 10.0,
388            )))
389            .unwrap();
390        let ema_after_second = 106.0 * k + 100.0 * (1.0 - k);
391        assert!((confirmed.unwrap().value - ema_after_second).abs() < 1e-9);
392
393        // Correcting the second bar's close must recompute EMA from the pre-second-bar state
394        // (EMA = 100.0 after the first bar), not from the already-advanced post-second-bar EMA.
395        let corrected = runner
396            .on_event(BarLifecycle::Correction(Bar::new(
397                1_060, 100.0, 109.0, 100.0, 109.0, 10.0,
398            )))
399            .unwrap();
400        let expected_corrected_ema = 109.0 * k + 100.0 * (1.0 - k);
401        assert!((corrected.unwrap().value - expected_corrected_ema).abs() < 1e-9);
402
403        // A double-application bug would instead apply 109.0 on top of the post-second-bar EMA
404        // (~104.0), producing a visibly different, larger value.
405        let buggy_double_applied = 109.0 * k + ema_after_second * (1.0 - k);
406        assert!((expected_corrected_ema - buggy_double_applied).abs() > 1e-6);
407    }
408}