1use super::*;
2
3pub fn claim_candidate(
4 checkpoint: &Checkpoint,
5 cycle_index: u64,
6 now_ms: u64,
7 claim_mode: ClaimMode,
8) -> CheckpointResult<bool> {
9 if cycle_index == 0 || cycle_index > MAX_WIRE_INTEGER {
10 return Err(CheckpointError::new(
11 "checkpoint_claim_invalid",
12 "claimed cycle must be positive and JSON-safe",
13 ));
14 }
15 if now_ms > MAX_WIRE_INTEGER {
16 return Err(CheckpointError::new(
17 "checkpoint_claim_invalid",
18 "now_ms is outside the JSON-safe range",
19 ));
20 }
21 if checkpoint.terminal_result.is_some() || checkpoint.status.is_terminal() {
22 return Ok(false);
23 }
24 if !matches!(
25 checkpoint.status,
26 CheckpointStatus::Running | CheckpointStatus::ReconciliationRequired
27 ) {
28 return Ok(false);
29 }
30 if checkpoint.cycle_index.checked_add(1) != Some(cycle_index) {
31 return Ok(false);
32 }
33 if (checkpoint.status == CheckpointStatus::ReconciliationRequired
34 || checkpoint.has_ambiguous_operation())
35 && claim_mode != ClaimMode::Recovery
36 {
37 return Ok(false);
38 }
39 if let Some(expiry) = checkpoint.lease_expires_at_ms {
40 if expiry > now_ms {
41 return Ok(false);
42 }
43 if claim_mode != ClaimMode::Recovery {
44 return Ok(false);
45 }
46 }
47 Ok(true)
48}
49
50pub fn apply_claim(
51 checkpoint: &mut Checkpoint,
52 cycle_index: u64,
53 claim_token: &str,
54 lease_expires_at_ms: u64,
55 claim_mode: ClaimMode,
56) -> CheckpointResult<()> {
57 if claim_token.trim().is_empty() || lease_expires_at_ms > MAX_WIRE_INTEGER {
58 return Err(CheckpointError::new(
59 "checkpoint_claim_invalid",
60 "claim token and lease must be non-empty and JSON-safe",
61 ));
62 }
63 checkpoint.revision = checkpoint
64 .revision
65 .checked_add(1)
66 .ok_or_else(|| CheckpointError::new("checkpoint_revision_overflow", "revision overflow"))?;
67 if claim_mode == ClaimMode::Recovery {
68 checkpoint.resume_attempt = checkpoint.resume_attempt.checked_add(1).ok_or_else(|| {
69 CheckpointError::new(
70 "checkpoint_resume_attempt_invalid",
71 "resume_attempt overflow",
72 )
73 })?;
74 }
75 checkpoint.status = CheckpointStatus::Running;
76 checkpoint.claim_token = Some(claim_token.to_string());
77 checkpoint.claimed_cycle = Some(cycle_index);
78 checkpoint.lease_expires_at_ms = Some(lease_expires_at_ms);
79 Ok(())
80}
81
82pub fn claim_matches(
83 current: &Checkpoint,
84 snapshot: &Checkpoint,
85 claim_token: &str,
86 expected_revision: u64,
87) -> bool {
88 current.revision == expected_revision
89 && snapshot.revision == expected_revision
90 && current.claim_token.as_deref() == Some(claim_token)
91 && current.claimed_cycle == snapshot.claimed_cycle
92 && current.checkpoint_key == snapshot.checkpoint_key
93 && current.terminal_result.is_none()
94 && checkpoint_definition_matches(current, snapshot)
95}
96
97pub fn checkpoint_definition_matches(current: &Checkpoint, snapshot: &Checkpoint) -> bool {
98 current.schema_version == snapshot.schema_version
99 && current.run_definition_schema == snapshot.run_definition_schema
100 && current.checkpoint_key == snapshot.checkpoint_key
101 && current.task_id == snapshot.task_id
102 && current.root_run_id == snapshot.root_run_id
103 && current.trace_id == snapshot.trace_id
104 && current.run_definition_digest == snapshot.run_definition_digest
105 && current.run_definition == snapshot.run_definition
106 && current.resume_attempt == snapshot.resume_attempt
107 && current.terminal_acknowledged == snapshot.terminal_acknowledged
108}
109
110pub fn prepare_progress(
111 current: &Checkpoint,
112 mut snapshot: Checkpoint,
113 claim_token: &str,
114 expected_revision: u64,
115) -> CheckpointResult<Option<Checkpoint>> {
116 if !claim_matches(current, &snapshot, claim_token, expected_revision) {
117 return Ok(None);
118 }
119 snapshot.claim_token = current.claim_token.clone();
120 snapshot.claimed_cycle = current.claimed_cycle;
121 snapshot.lease_expires_at_ms = current.lease_expires_at_ms;
122 snapshot.revision = expected_revision
123 .checked_add(1)
124 .ok_or_else(|| CheckpointError::new("checkpoint_revision_overflow", "revision overflow"))?;
125 snapshot.validate()?;
126 Ok(Some(snapshot))
127}
128
129pub fn prepare_suspend(
130 current: &Checkpoint,
131 mut snapshot: Checkpoint,
132 claim_token: &str,
133 expected_revision: u64,
134) -> CheckpointResult<Option<Checkpoint>> {
135 if !claim_matches(current, &snapshot, claim_token, expected_revision)
136 || !snapshot.has_ambiguous_operation()
137 {
138 return Ok(None);
139 }
140 snapshot.status = CheckpointStatus::ReconciliationRequired;
141 snapshot.claim_token = None;
142 snapshot.claimed_cycle = None;
143 snapshot.lease_expires_at_ms = None;
144 snapshot.revision = expected_revision
145 .checked_add(1)
146 .ok_or_else(|| CheckpointError::new("checkpoint_revision_overflow", "revision overflow"))?;
147 snapshot.validate()?;
148 Ok(Some(snapshot))
149}
150
151pub fn prepare_commit(
152 current: &Checkpoint,
153 mut snapshot: Checkpoint,
154 claim_token: &str,
155 expected_revision: u64,
156) -> CheckpointResult<Option<Checkpoint>> {
157 if !claim_matches(current, &snapshot, claim_token, expected_revision) {
158 return Ok(None);
159 }
160 let Some(claimed_cycle) = current.claimed_cycle else {
161 return Ok(None);
162 };
163 if snapshot.cycle_index != claimed_cycle {
164 return Ok(None);
165 }
166 validate_model_journal_accounting(&snapshot)?;
167 snapshot
168 .event_outbox
169 .retain(|entry| entry.state == "pending");
170 snapshot.model_call_journal.clear();
171 snapshot.tool_journal.clear();
172 snapshot.claim_token = None;
173 snapshot.claimed_cycle = None;
174 snapshot.lease_expires_at_ms = None;
175 snapshot.revision = expected_revision
176 .checked_add(1)
177 .ok_or_else(|| CheckpointError::new("checkpoint_revision_overflow", "revision overflow"))?;
178 snapshot.validate()?;
179 Ok(Some(snapshot))
180}
181
182pub fn prepare_finalize(
183 current: &Checkpoint,
184 snapshot: Checkpoint,
185 expected_revision: u64,
186) -> CheckpointResult<Option<Checkpoint>> {
187 if current.revision != expected_revision
188 || snapshot.revision != expected_revision
189 || !checkpoint_definition_matches(current, &snapshot)
190 || current.claim_token.is_some()
191 || current.terminal_result.is_some()
192 {
193 return Ok(None);
194 }
195 prepare_terminal_snapshot(snapshot, expected_revision).map(Some)
196}
197
198pub fn prepare_finalize_claimed(
199 current: &Checkpoint,
200 snapshot: Checkpoint,
201 claim_token: &str,
202 expected_revision: u64,
203) -> CheckpointResult<Option<Checkpoint>> {
204 if !claim_matches(current, &snapshot, claim_token, expected_revision) {
205 return Ok(None);
206 }
207 prepare_terminal_snapshot(snapshot, expected_revision).map(Some)
208}
209
210fn prepare_terminal_snapshot(
211 mut snapshot: Checkpoint,
212 expected_revision: u64,
213) -> CheckpointResult<Checkpoint> {
214 let Some(terminal_result) = snapshot.terminal_result.as_ref() else {
215 return Err(CheckpointError::new(
216 "checkpoint_terminal_result_required",
217 "finalize requires terminal_result",
218 ));
219 };
220 if !snapshot.status.is_terminal() {
221 return Err(CheckpointError::new(
222 "checkpoint_status_invalid",
223 "finalize requires a terminal status",
224 ));
225 }
226 let operator_abort = snapshot.is_operator_abort_terminal();
227 validate_model_journal_accounting(&snapshot)?;
228 snapshot
229 .event_outbox
230 .retain(|entry| entry.state == "pending");
231 if !operator_abort {
232 snapshot.model_call_journal.clear();
233 snapshot.tool_journal.clear();
234 }
235 snapshot.claim_token = None;
236 snapshot.claimed_cycle = None;
237 snapshot.lease_expires_at_ms = None;
238 snapshot.revision = expected_revision
239 .checked_add(1)
240 .ok_or_else(|| CheckpointError::new("checkpoint_revision_overflow", "revision overflow"))?;
241 let _ = terminal_result;
242 snapshot.validate()?;
243 Ok(snapshot)
244}
245
246pub fn prepare_event_delivery(
247 current: &Checkpoint,
248 claim_token: Option<&str>,
249 expected_revision: u64,
250 event_id: &str,
251 payload_digest: &str,
252 cursor: EventCursor,
253) -> CheckpointResult<Option<Checkpoint>> {
254 if event_id.trim().is_empty() {
255 return Err(CheckpointError::new(
256 "checkpoint_event_outbox_invalid",
257 "event_id must be non-empty",
258 ));
259 }
260 validate_sha256(payload_digest, "event_outbox.payload_digest")?;
261 cursor.validate()?;
262 if cursor.last_event_id.as_deref() != Some(event_id) {
263 return Err(CheckpointError::new(
264 "checkpoint_event_cursor_invalid",
265 "event cursor last_event_id must match the delivered event",
266 ));
267 }
268 if current.revision != expected_revision || current.claim_token.as_deref() != claim_token {
269 return Ok(None);
270 }
271
272 let matching = current
273 .event_outbox
274 .iter()
275 .enumerate()
276 .filter(|(_, entry)| entry.event_id == event_id)
277 .collect::<Vec<_>>();
278 if matching.len() != 1 {
279 return Ok(None);
280 }
281 let (index, entry) = matching[0];
282 if entry.state != "pending" || entry.payload_digest != payload_digest {
283 return Ok(None);
284 }
285
286 let cursor_value = serde_json::to_value(&cursor).map_err(|error| {
287 CheckpointError::new("checkpoint_event_cursor_invalid", error.to_string())
288 })?;
289 let mut snapshot = current.clone();
290 snapshot.event_outbox[index].state = "delivered".to_string();
291 snapshot.event_outbox[index].cursor = Some(cursor_value);
292 snapshot.event_cursor = Some(cursor);
293 snapshot.revision = expected_revision
294 .checked_add(1)
295 .ok_or_else(|| CheckpointError::new("checkpoint_revision_overflow", "revision overflow"))?;
296 snapshot.validate()?;
297 Ok(Some(snapshot))
298}
299
300pub fn prepare_ack(
301 current: &Checkpoint,
302 expected_revision: u64,
303) -> CheckpointResult<Option<Checkpoint>> {
304 if current.revision != expected_revision
305 || current.terminal_result.is_none()
306 || current.claim_token.is_some()
307 || current.terminal_acknowledged
308 {
309 return Ok(None);
310 }
311 let mut snapshot = current.clone();
312 snapshot.terminal_acknowledged = true;
313 snapshot.revision = expected_revision
314 .checked_add(1)
315 .ok_or_else(|| CheckpointError::new("checkpoint_revision_overflow", "revision overflow"))?;
316 snapshot.validate()?;
317 Ok(Some(snapshot))
318}