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}