Skip to main content

vv_agent/runtime/state_v2/
transitions.rs

1use super::*;
2
3pub fn claim_candidate(
4    checkpoint: &CheckpointV2,
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 CheckpointV2,
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: &CheckpointV2,
84    snapshot: &CheckpointV2,
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: &CheckpointV2, snapshot: &CheckpointV2) -> 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: &CheckpointV2,
112    mut snapshot: CheckpointV2,
113    claim_token: &str,
114    expected_revision: u64,
115) -> CheckpointResult<Option<CheckpointV2>> {
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: &CheckpointV2,
131    mut snapshot: CheckpointV2,
132    claim_token: &str,
133    expected_revision: u64,
134) -> CheckpointResult<Option<CheckpointV2>> {
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: &CheckpointV2,
153    mut snapshot: CheckpointV2,
154    claim_token: &str,
155    expected_revision: u64,
156) -> CheckpointResult<Option<CheckpointV2>> {
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    snapshot.model_call_journal.clear();
167    snapshot.tool_journal.clear();
168    snapshot.claim_token = None;
169    snapshot.claimed_cycle = None;
170    snapshot.lease_expires_at_ms = None;
171    snapshot.revision = expected_revision
172        .checked_add(1)
173        .ok_or_else(|| CheckpointError::new("checkpoint_revision_overflow", "revision overflow"))?;
174    snapshot.validate()?;
175    Ok(Some(snapshot))
176}
177
178pub fn prepare_finalize(
179    current: &CheckpointV2,
180    snapshot: CheckpointV2,
181    expected_revision: u64,
182) -> CheckpointResult<Option<CheckpointV2>> {
183    if current.revision != expected_revision
184        || snapshot.revision != expected_revision
185        || !checkpoint_definition_matches(current, &snapshot)
186        || current.claim_token.is_some()
187        || current.terminal_result.is_some()
188    {
189        return Ok(None);
190    }
191    prepare_terminal_snapshot(snapshot, expected_revision).map(Some)
192}
193
194pub fn prepare_finalize_claimed(
195    current: &CheckpointV2,
196    snapshot: CheckpointV2,
197    claim_token: &str,
198    expected_revision: u64,
199) -> CheckpointResult<Option<CheckpointV2>> {
200    if !claim_matches(current, &snapshot, claim_token, expected_revision) {
201        return Ok(None);
202    }
203    prepare_terminal_snapshot(snapshot, expected_revision).map(Some)
204}
205
206fn prepare_terminal_snapshot(
207    mut snapshot: CheckpointV2,
208    expected_revision: u64,
209) -> CheckpointResult<CheckpointV2> {
210    let Some(terminal_result) = snapshot.terminal_result.as_ref() else {
211        return Err(CheckpointError::new(
212            "checkpoint_terminal_result_required",
213            "finalize requires terminal_result",
214        ));
215    };
216    if !snapshot.status.is_terminal() {
217        return Err(CheckpointError::new(
218            "checkpoint_status_invalid",
219            "finalize requires a terminal status",
220        ));
221    }
222    let operator_abort = snapshot.is_operator_abort_terminal();
223    if !operator_abort {
224        snapshot.model_call_journal.clear();
225        snapshot.tool_journal.clear();
226    }
227    snapshot.claim_token = None;
228    snapshot.claimed_cycle = None;
229    snapshot.lease_expires_at_ms = None;
230    snapshot.revision = expected_revision
231        .checked_add(1)
232        .ok_or_else(|| CheckpointError::new("checkpoint_revision_overflow", "revision overflow"))?;
233    let _ = terminal_result;
234    snapshot.validate()?;
235    Ok(snapshot)
236}
237
238pub fn prepare_event_delivery(
239    current: &CheckpointV2,
240    claim_token: Option<&str>,
241    expected_revision: u64,
242    event_id: &str,
243    payload_digest: &str,
244    cursor: EventCursor,
245) -> CheckpointResult<Option<CheckpointV2>> {
246    if event_id.trim().is_empty() {
247        return Err(CheckpointError::new(
248            "checkpoint_event_outbox_invalid",
249            "event_id must be non-empty",
250        ));
251    }
252    validate_sha256(payload_digest, "event_outbox.payload_digest")?;
253    cursor.validate()?;
254    if cursor.last_event_id.as_deref() != Some(event_id) {
255        return Err(CheckpointError::new(
256            "checkpoint_event_cursor_invalid",
257            "event cursor last_event_id must match the delivered event",
258        ));
259    }
260    if current.revision != expected_revision || current.claim_token.as_deref() != claim_token {
261        return Ok(None);
262    }
263
264    let matching = current
265        .event_outbox
266        .iter()
267        .enumerate()
268        .filter(|(_, entry)| entry.event_id == event_id)
269        .collect::<Vec<_>>();
270    if matching.len() != 1 {
271        return Ok(None);
272    }
273    let (index, entry) = matching[0];
274    if entry.state != "pending" || entry.payload_digest != payload_digest {
275        return Ok(None);
276    }
277
278    let cursor_value = serde_json::to_value(&cursor).map_err(|error| {
279        CheckpointError::new("checkpoint_event_cursor_invalid", error.to_string())
280    })?;
281    let mut snapshot = current.clone();
282    snapshot.event_outbox[index].state = "delivered".to_string();
283    snapshot.event_outbox[index].cursor = Some(cursor_value);
284    snapshot.event_cursor = Some(cursor);
285    snapshot.revision = expected_revision
286        .checked_add(1)
287        .ok_or_else(|| CheckpointError::new("checkpoint_revision_overflow", "revision overflow"))?;
288    snapshot.validate()?;
289    Ok(Some(snapshot))
290}
291
292pub fn prepare_ack(
293    current: &CheckpointV2,
294    expected_revision: u64,
295) -> CheckpointResult<Option<CheckpointV2>> {
296    if current.revision != expected_revision
297        || current.terminal_result.is_none()
298        || current.claim_token.is_some()
299        || current.terminal_acknowledged
300    {
301        return Ok(None);
302    }
303    let mut snapshot = current.clone();
304    snapshot.terminal_acknowledged = true;
305    snapshot.revision = expected_revision
306        .checked_add(1)
307        .ok_or_else(|| CheckpointError::new("checkpoint_revision_overflow", "revision overflow"))?;
308    snapshot.validate()?;
309    Ok(Some(snapshot))
310}