Skip to main content

kestrel_chartkit/
checkpoint.rs

1//! Versioned state snapshots for long-running engines.
2//!
3//! Any `Clone` state (an `Indicator`, an `engine::*` tracker, a composition graph) can be captured
4//! into a [`Checkpoint`] and restored later, in-process, to deterministically save/rewind
5//! execution. When `S` also implements `serde::Serialize`/`Deserialize` (behind the `serde`
6//! feature), the checkpoint itself becomes serializable for cross-process persistence — no
7//! separate persistence API is required.
8
9#[cfg(feature = "serde")]
10use serde::{Deserialize, Serialize};
11
12/// A versioned, timestamped snapshot of some `Clone`-able engine state `S`.
13#[derive(Debug, Clone, PartialEq)]
14#[cfg_attr(feature = "serde", derive(Serialize, Deserialize))]
15pub struct Checkpoint<S: Clone> {
16    /// Monotonically increasing checkpoint version, incremented by the caller on each capture.
17    /// Lets consumers detect stale/out-of-order restores.
18    pub version: u32,
19    /// The bar timestamp the state reflects (i.e. state after processing this bar).
20    pub timestamp: i64,
21    pub state: S,
22}
23
24impl<S: Clone> Checkpoint<S> {
25    /// Captures a snapshot of `state` at `timestamp` tagged with `version`.
26    pub fn capture(state: &S, timestamp: i64, version: u32) -> Self {
27        Self {
28            version,
29            timestamp,
30            state: state.clone(),
31        }
32    }
33
34    /// Overwrites `target` with this checkpoint's state.
35    pub fn restore_into(&self, target: &mut S) {
36        *target = self.state.clone();
37    }
38}
39
40/// Keeps the single most recent [`Checkpoint`] for `S`, auto-incrementing the version on every
41/// [`CheckpointStore::save`]. Suited for a "last confirmed state" rewind point on a long-running
42/// engine.
43#[derive(Debug, Clone, Default)]
44pub struct CheckpointStore<S: Clone> {
45    latest: Option<Checkpoint<S>>,
46    next_version: u32,
47}
48
49impl<S: Clone> CheckpointStore<S> {
50    pub fn new() -> Self {
51        Self {
52            latest: None,
53            next_version: 0,
54        }
55    }
56
57    /// Captures `state` as the new latest checkpoint, replacing any prior one.
58    pub fn save(&mut self, state: &S, timestamp: i64) -> u32 {
59        let version = self.next_version;
60        self.next_version += 1;
61        self.latest = Some(Checkpoint::capture(state, timestamp, version));
62        version
63    }
64
65    pub fn latest(&self) -> Option<&Checkpoint<S>> {
66        self.latest.as_ref()
67    }
68
69    /// Restores `target` from the latest checkpoint, if one exists.
70    pub fn restore_into(&self, target: &mut S) -> bool {
71        match &self.latest {
72            Some(checkpoint) => {
73                checkpoint.restore_into(target);
74                true
75            }
76            None => false,
77        }
78    }
79}
80
81#[cfg(test)]
82mod tests {
83    use super::*;
84
85    #[test]
86    fn test_checkpoint_capture_and_restore() {
87        let state = vec![1, 2, 3];
88        let checkpoint = Checkpoint::capture(&state, 1_000, 0);
89        assert_eq!(checkpoint.version, 0);
90        assert_eq!(checkpoint.timestamp, 1_000);
91
92        let mut target = vec![9, 9];
93        checkpoint.restore_into(&mut target);
94        assert_eq!(target, state);
95    }
96
97    #[test]
98    fn test_checkpoint_store_versions_increment() {
99        let mut store: CheckpointStore<i32> = CheckpointStore::new();
100        assert!(store.latest().is_none());
101
102        let v0 = store.save(&10, 1_000);
103        let v1 = store.save(&20, 1_060);
104        assert_eq!(v0, 0);
105        assert_eq!(v1, 1);
106        assert_eq!(store.latest().unwrap().state, 20);
107
108        let mut target = 0;
109        assert!(store.restore_into(&mut target));
110        assert_eq!(target, 20);
111    }
112}