Skip to main content

vv_agent/runtime/state/
transitions.rs

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}