Skip to main content

vv_agent/runtime/stores/
memory_v2.rs

1//! In-memory checkpoint v2 store.
2
3use std::collections::BTreeMap;
4use std::sync::{Arc, Mutex};
5
6use crate::checkpoint::{CheckpointError, CheckpointResult, ClaimMode, EventCursor};
7use crate::runtime::state_v2::{
8    apply_claim, claim_candidate, prepare_ack, prepare_commit, prepare_event_delivery,
9    prepare_finalize, prepare_finalize_claimed, prepare_progress, prepare_suspend,
10    CheckpointStoreV2, CheckpointV2,
11};
12
13#[derive(Debug, Clone, Default)]
14pub struct InMemoryCheckpointStoreV2 {
15    checkpoints: Arc<Mutex<BTreeMap<String, CheckpointV2>>>,
16}
17
18impl InMemoryCheckpointStoreV2 {
19    pub fn new() -> Self {
20        Self::default()
21    }
22
23    pub fn save_checkpoint_v2(&self, checkpoint: CheckpointV2) -> CheckpointResult<()> {
24        checkpoint.validate()?;
25        let mut checkpoints = self.lock()?;
26        checkpoints.insert(checkpoint.checkpoint_key.clone(), checkpoint);
27        Ok(())
28    }
29
30    fn lock(&self) -> CheckpointResult<std::sync::MutexGuard<'_, BTreeMap<String, CheckpointV2>>> {
31        self.checkpoints.lock().map_err(|_| {
32            CheckpointError::new(
33                "checkpoint_store_lock_poisoned",
34                "checkpoint store lock poisoned",
35            )
36        })
37    }
38}
39
40impl CheckpointStoreV2 for InMemoryCheckpointStoreV2 {
41    fn create_checkpoint_v2(&self, checkpoint: CheckpointV2) -> CheckpointResult<bool> {
42        checkpoint.validate()?;
43        let mut checkpoints = self.lock()?;
44        if checkpoints.contains_key(&checkpoint.checkpoint_key) {
45            return Ok(false);
46        }
47        checkpoints.insert(checkpoint.checkpoint_key.clone(), checkpoint);
48        Ok(true)
49    }
50
51    fn load_checkpoint_v2(&self, checkpoint_key: &str) -> CheckpointResult<Option<CheckpointV2>> {
52        let checkpoints = self.lock()?;
53        Ok(checkpoints.get(checkpoint_key).cloned())
54    }
55
56    fn claim_checkpoint_v2(
57        &self,
58        checkpoint_key: &str,
59        cycle_index: u64,
60        claim_token: &str,
61        lease_expires_at_ms: u64,
62        now_ms: u64,
63        claim_mode: ClaimMode,
64    ) -> CheckpointResult<Option<CheckpointV2>> {
65        if claim_token.trim().is_empty() || lease_expires_at_ms <= now_ms {
66            return Err(CheckpointError::new(
67                "checkpoint_claim_invalid",
68                "claim token must be non-empty and lease must be in the future",
69            ));
70        }
71        let mut checkpoints = self.lock()?;
72        let Some(current) = checkpoints.get(checkpoint_key).cloned() else {
73            return Ok(None);
74        };
75        if !claim_candidate(&current, cycle_index, now_ms, claim_mode)? {
76            return Ok(None);
77        }
78        let mut claimed = current;
79        apply_claim(
80            &mut claimed,
81            cycle_index,
82            claim_token,
83            lease_expires_at_ms,
84            claim_mode,
85        )?;
86        claimed.validate()?;
87        checkpoints.insert(checkpoint_key.to_string(), claimed.clone());
88        Ok(Some(claimed))
89    }
90
91    fn progress_checkpoint_v2(
92        &self,
93        checkpoint: CheckpointV2,
94        claim_token: &str,
95        expected_revision: u64,
96    ) -> CheckpointResult<bool> {
97        let mut checkpoints = self.lock()?;
98        let Some(current) = checkpoints.get(&checkpoint.checkpoint_key).cloned() else {
99            return Ok(false);
100        };
101        let Some(updated) = prepare_progress(&current, checkpoint, claim_token, expected_revision)?
102        else {
103            return Ok(false);
104        };
105        checkpoints.insert(updated.checkpoint_key.clone(), updated);
106        Ok(true)
107    }
108
109    fn suspend_checkpoint_v2(
110        &self,
111        checkpoint: CheckpointV2,
112        claim_token: &str,
113        expected_revision: u64,
114    ) -> CheckpointResult<bool> {
115        let mut checkpoints = self.lock()?;
116        let Some(current) = checkpoints.get(&checkpoint.checkpoint_key).cloned() else {
117            return Ok(false);
118        };
119        let Some(updated) = prepare_suspend(&current, checkpoint, claim_token, expected_revision)?
120        else {
121            return Ok(false);
122        };
123        checkpoints.insert(updated.checkpoint_key.clone(), updated);
124        Ok(true)
125    }
126
127    fn commit_checkpoint_v2(
128        &self,
129        checkpoint: CheckpointV2,
130        claim_token: &str,
131        expected_revision: u64,
132    ) -> CheckpointResult<bool> {
133        let mut checkpoints = self.lock()?;
134        let Some(current) = checkpoints.get(&checkpoint.checkpoint_key).cloned() else {
135            return Ok(false);
136        };
137        let Some(updated) = prepare_commit(&current, checkpoint, claim_token, expected_revision)?
138        else {
139            return Ok(false);
140        };
141        checkpoints.insert(updated.checkpoint_key.clone(), updated);
142        Ok(true)
143    }
144
145    fn finalize_checkpoint_v2(
146        &self,
147        checkpoint: CheckpointV2,
148        expected_revision: u64,
149    ) -> CheckpointResult<bool> {
150        let mut checkpoints = self.lock()?;
151        let Some(current) = checkpoints.get(&checkpoint.checkpoint_key).cloned() else {
152            return Ok(false);
153        };
154        let Some(updated) = prepare_finalize(&current, checkpoint, expected_revision)? else {
155            return Ok(false);
156        };
157        checkpoints.insert(updated.checkpoint_key.clone(), updated);
158        Ok(true)
159    }
160
161    fn finalize_claimed_v2(
162        &self,
163        checkpoint: CheckpointV2,
164        claim_token: &str,
165        expected_revision: u64,
166    ) -> CheckpointResult<bool> {
167        let mut checkpoints = self.lock()?;
168        let Some(current) = checkpoints.get(&checkpoint.checkpoint_key).cloned() else {
169            return Ok(false);
170        };
171        let Some(updated) =
172            prepare_finalize_claimed(&current, checkpoint, claim_token, expected_revision)?
173        else {
174            return Ok(false);
175        };
176        checkpoints.insert(updated.checkpoint_key.clone(), updated);
177        Ok(true)
178    }
179
180    fn renew_checkpoint_claim_v2(
181        &self,
182        checkpoint_key: &str,
183        claim_token: &str,
184        lease_expires_at_ms: u64,
185        now_ms: u64,
186    ) -> CheckpointResult<bool> {
187        if claim_token.trim().is_empty() || lease_expires_at_ms <= now_ms {
188            return Err(CheckpointError::new(
189                "checkpoint_claim_invalid",
190                "claim token must be non-empty and lease must be in the future",
191            ));
192        }
193        let mut checkpoints = self.lock()?;
194        let Some(checkpoint) = checkpoints.get_mut(checkpoint_key) else {
195            return Ok(false);
196        };
197        if checkpoint.claim_token.as_deref() != Some(claim_token)
198            || checkpoint
199                .lease_expires_at_ms
200                .is_none_or(|expiry| expiry <= now_ms)
201        {
202            return Ok(false);
203        }
204        checkpoint.lease_expires_at_ms = Some(lease_expires_at_ms);
205        Ok(true)
206    }
207
208    fn acknowledge_terminal_v2(
209        &self,
210        checkpoint_key: &str,
211        expected_revision: u64,
212    ) -> CheckpointResult<bool> {
213        let mut checkpoints = self.lock()?;
214        let Some(current) = checkpoints.get(checkpoint_key).cloned() else {
215            return Ok(false);
216        };
217        let Some(updated) = prepare_ack(&current, expected_revision)? else {
218            return Ok(false);
219        };
220        checkpoints.insert(checkpoint_key.to_string(), updated);
221        Ok(true)
222    }
223
224    fn record_event_delivery_v2(
225        &self,
226        checkpoint_key: &str,
227        claim_token: Option<&str>,
228        expected_revision: u64,
229        event_id: &str,
230        payload_digest: &str,
231        cursor: EventCursor,
232    ) -> CheckpointResult<bool> {
233        let mut checkpoints = self.lock()?;
234        let Some(current) = checkpoints.get(checkpoint_key).cloned() else {
235            return Ok(false);
236        };
237        let Some(updated) = prepare_event_delivery(
238            &current,
239            claim_token,
240            expected_revision,
241            event_id,
242            payload_digest,
243            cursor,
244        )?
245        else {
246            return Ok(false);
247        };
248        checkpoints.insert(checkpoint_key.to_string(), updated);
249        Ok(true)
250    }
251
252    fn delete_checkpoint_v2(&self, checkpoint_key: &str) -> CheckpointResult<()> {
253        self.lock()?.remove(checkpoint_key);
254        Ok(())
255    }
256
257    fn list_checkpoints_v2(&self) -> CheckpointResult<Vec<String>> {
258        Ok(self.lock()?.keys().cloned().collect())
259    }
260}
261
262pub type InMemoryStateStoreV2 = InMemoryCheckpointStoreV2;