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}