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//! (Pine's `barstate.isconfirmed == false`) versus one that has closed, and no way to reprocess a
6//! bar without double-counting it.
7//!
8//! [`BarLifecycle`] names the four events a feed can emit for a given timestamp. [`LifecycleRunner`]
9//! wraps any `Indicator + Clone` and makes [`BarLifecycle::Update`]/[`BarLifecycle::Correction`]
10//! safe to call repeatedly: it always recomputes from the last *confirmed* checkpoint rather than
11//! mutating forward, so re-delivering the same revised bar is idempotent and never double-applies
12//! a still-forming bar's data.
13//!
14//! Requiring `Clone` is a deliberate, honest scope limit: it is the only mechanism that works
15//! generically across arbitrary indicator internals without adding a second, hand-written
16//! snapshot/restore method to every `Indicator` impl. Most indicator structs in this crate do not
17//! yet derive `Clone`; adding it is a mechanical, low-risk follow-up left to indicator authors as
18//! they adopt `LifecycleRunner`, not a blanket change bundled into this module.
19
20use crate::checkpoint::CheckpointStore;
21use crate::indicator::{Indicator, IndicatorOutput};
22use crate::model::Bar;
23
24/// A bar-lifecycle event for a single timestamp.
25#[derive(Debug, Clone, PartialEq)]
26pub enum BarLifecycle {
27    /// The first tick of a new bar/timestamp.
28    Open(Bar),
29    /// A subsequent, still-unconfirmed tick for the same timestamp as the last `Open`/`Update`
30    /// (Pine's `barstate.isconfirmed == false`).
31    Update(Bar),
32    /// The bar for this timestamp is final and will not be revised again.
33    Confirmed(Bar),
34    /// A feed correction to the most recently confirmed bar (e.g. an exchange restatement).
35    /// Only single-level undo is supported: correcting a bar older than the last confirmed one
36    /// is not representable by [`LifecycleRunner`], which keeps just one checkpoint.
37    Correction(Bar),
38}
39
40impl BarLifecycle {
41    pub fn bar(&self) -> &Bar {
42        match self {
43            BarLifecycle::Open(b)
44            | BarLifecycle::Update(b)
45            | BarLifecycle::Confirmed(b)
46            | BarLifecycle::Correction(b) => b,
47        }
48    }
49}
50
51/// Wraps an `Indicator + Clone` so [`BarLifecycle::Open`]/[`BarLifecycle::Update`] events can be
52/// replayed any number of times without permanently mutating state, and only
53/// [`BarLifecycle::Confirmed`]/[`BarLifecycle::Correction`] commit a new checkpoint.
54pub struct LifecycleRunner<I: Indicator + Clone> {
55    live: I,
56    checkpoints: CheckpointStore<I>,
57}
58
59impl<I: Indicator + Clone> LifecycleRunner<I> {
60    pub fn new(indicator: I) -> Self {
61        let checkpoints = CheckpointStore::new();
62        Self {
63            live: indicator,
64            checkpoints,
65        }
66    }
67
68    /// The indicator state as of the last confirmed/corrected bar (never a tentative
69    /// open/update). `None` before the first confirmation.
70    pub fn confirmed_indicator(&self) -> Option<&I> {
71        self.checkpoints.latest().map(|c| &c.state)
72    }
73
74    /// Processes one lifecycle event and returns the resulting output.
75    ///
76    /// `Open`/`Update` always start from the last confirmed checkpoint (or a freshly reset
77    /// indicator if none exists yet) before applying `bar`, so repeated calls for the same
78    /// still-forming timestamp are idempotent and never stack on top of a prior tentative
79    /// application. `Confirmed`/`Correction` do the same, then persist the resulting state as the
80    /// new checkpoint.
81    pub fn on_event(&mut self, event: BarLifecycle) -> Option<IndicatorOutput> {
82        self.rewind_to_last_confirmed();
83        let output = self.live.on_bar(event.bar());
84
85        if matches!(
86            event,
87            BarLifecycle::Confirmed(_) | BarLifecycle::Correction(_)
88        ) {
89            self.checkpoints.save(&self.live, event.bar().timestamp);
90        }
91
92        output
93    }
94
95    fn rewind_to_last_confirmed(&mut self) {
96        if !self.checkpoints.restore_into(&mut self.live) {
97            self.live.reset();
98        }
99    }
100}
101
102#[cfg(test)]
103mod tests {
104    use super::*;
105
106    #[test]
107    fn test_repeated_update_is_idempotent() {
108        // SmaEngine does not derive Clone (see module docs), so this test uses a minimal
109        // hand-rolled Clone indicator instead of pulling SmaEngine into the Clone surface.
110        #[derive(Clone)]
111        struct SumEngine {
112            sum: f64,
113        }
114        impl Indicator for SumEngine {
115            fn name(&self) -> &str {
116                "sum"
117            }
118            fn reset(&mut self) {
119                self.sum = 0.0;
120            }
121            fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
122                self.sum += bar.close;
123                Some(IndicatorOutput::new(self.sum))
124            }
125        }
126
127        let mut runner = LifecycleRunner::new(SumEngine { sum: 0.0 });
128        let bar = Bar::new(1_000, 100.0, 101.0, 99.0, 100.0, 10.0);
129
130        let first_update = runner.on_event(BarLifecycle::Open(bar.clone()));
131        let second_update = runner.on_event(BarLifecycle::Update(bar.clone()));
132        // Re-delivering the same still-forming bar must not accumulate.
133        assert_eq!(first_update.unwrap().value, second_update.unwrap().value);
134
135        let confirmed = runner.on_event(BarLifecycle::Confirmed(bar.clone()));
136        assert_eq!(confirmed.unwrap().value, 100.0);
137        assert_eq!(runner.confirmed_indicator().unwrap().sum, 100.0);
138
139        // A later bar accumulates on top of the confirmed checkpoint, not the discarded update.
140        let next_bar = Bar::new(1_060, 100.0, 102.0, 100.0, 101.0, 10.0);
141        let next = runner.on_event(BarLifecycle::Confirmed(next_bar));
142        assert_eq!(next.unwrap().value, 201.0);
143    }
144
145    #[test]
146    fn test_correction_replaces_last_confirmed_bar() {
147        #[derive(Clone)]
148        struct LastCloseEngine {
149            last_close: f64,
150        }
151        impl Indicator for LastCloseEngine {
152            fn name(&self) -> &str {
153                "last_close"
154            }
155            fn reset(&mut self) {
156                self.last_close = 0.0;
157            }
158            fn on_bar(&mut self, bar: &Bar) -> Option<IndicatorOutput> {
159                self.last_close = bar.close;
160                Some(IndicatorOutput::new(self.last_close))
161            }
162        }
163
164        let mut runner = LifecycleRunner::new(LastCloseEngine { last_close: 0.0 });
165        let bar = Bar::new(1_000, 100.0, 101.0, 99.0, 100.0, 10.0);
166        runner.on_event(BarLifecycle::Confirmed(bar));
167
168        let corrected_bar = Bar::new(1_000, 100.0, 101.0, 99.0, 103.5, 10.0);
169        let corrected = runner.on_event(BarLifecycle::Correction(corrected_bar));
170        assert_eq!(corrected.unwrap().value, 103.5);
171        assert_eq!(runner.confirmed_indicator().unwrap().last_close, 103.5);
172    }
173}