1use 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(¤t, 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(¤t, 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(¤t, 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(¤t, 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(¤t, 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(¤t, 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(¤t, 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 ¤t,
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;