Skip to main content

vtcode_memory/
matrix.rs

1//! Deterministic matrix projection and atomic local resource admission.
2use std::collections::{BTreeMap, BTreeSet};
3use std::path::{Component, Path};
4use thiserror::Error;
5use vtcode_exec_events::{ThreadEvent, matrix::*};
6
7#[derive(Debug, Error)]
8#[error("{0}")]
9/// Specification, transition, or recovery validation failure.
10pub struct MatrixError(pub String);
11/// Result of a matrix state transition.
12pub type MatrixResult<T> = Result<T, MatrixError>;
13
14fn require(condition: bool, message: impl Into<String>) -> MatrixResult<()> {
15    if condition {
16        Ok(())
17    } else {
18        Err(MatrixError(message.into()))
19    }
20}
21
22/// Validate source-independent specification structure before persistence or replay.
23pub fn validate_spec(spec: &MatrixSpec) -> MatrixResult<()> {
24    require(!spec.id.trim().is_empty() && !spec.tasks.is_empty(), "matrix needs an ID and tasks")?;
25    require(
26        spec.resources.iter().all(|(id, n)| !id.trim().is_empty() && *n > 0),
27        "resource capacities must be positive",
28    )?;
29    let ids: BTreeSet<_> = spec.tasks.iter().map(|task| task.id.as_str()).collect();
30    require(ids.len() == spec.tasks.len(), "duplicate task IDs")?;
31    for task in &spec.tasks {
32        require(
33            !task.id.trim().is_empty() && !task.instructions.trim().is_empty(),
34            "task needs an ID and instructions",
35        )?;
36        require(task.timeout_secs > 0, "task timeout must be positive")?;
37        require(
38            !task.checks.is_empty() && task.checks.iter().all(|check| !check.trim().is_empty()),
39            "task needs nonempty verification commands",
40        )?;
41        validate_relative(&task.workspace)?;
42        for input in &task.inputs {
43            validate_relative(input)?;
44        }
45        require(task.dependencies.iter().all(|id| ids.contains(id.as_str())), "unknown dependency")?;
46        require(
47            task.dependencies.iter().collect::<BTreeSet<_>>().len() == task.dependencies.len(),
48            "duplicate dependency",
49        )?;
50        for (resource, count) in &task.resources {
51            require(
52                *count > 0 && spec.resources.get(resource).is_some_and(|capacity| count <= capacity),
53                format!("invalid resource requirement: {resource}"),
54            )?;
55        }
56    }
57    let mut visited = BTreeSet::new();
58    loop {
59        let previous = visited.len();
60        for task in &spec.tasks {
61            if task.dependencies.iter().all(|id| visited.contains(id.as_str())) {
62                visited.insert(task.id.as_str());
63            }
64        }
65        if visited.len() == spec.tasks.len() {
66            return Ok(());
67        }
68        require(visited.len() > previous, "dependency cycle")?;
69    }
70}
71
72fn validate_relative(path: &str) -> MatrixResult<()> {
73    require(
74        !path.is_empty()
75            && !Path::new(path).is_absolute()
76            && Path::new(path)
77                .components()
78                .all(|component| matches!(component, Component::Normal(_) | Component::CurDir)),
79        "workspace and input paths must stay relative to the workspace",
80    )
81}
82
83/// Reject workspace symlink escapes and unavailable declared inputs before dispatch.
84pub fn validate_workspace(spec: &MatrixSpec, root: &Path) -> MatrixResult<()> {
85    let canonical_root =
86        vtcode_commons::canonicalize(root).map_err(|error| MatrixError(format!("resolve workspace: {error}")))?;
87    for task in &spec.tasks {
88        let workspace = vtcode_commons::canonicalize(root.join(&task.workspace))
89            .map_err(|error| MatrixError(format!("resolve task workspace {}: {error}", task.id)))?;
90        require(
91            workspace.starts_with(&canonical_root) && workspace.is_dir(),
92            "task workspace escapes root or is not a directory",
93        )?;
94        for input in &task.inputs {
95            let path = vtcode_commons::canonicalize(workspace.join(input))
96                .map_err(|error| MatrixError(format!("resolve input {input}: {error}")))?;
97            require(
98                path.starts_with(&canonical_root) && path.is_file(),
99                "declared input escapes root or is not a file",
100            )?;
101        }
102    }
103    Ok(())
104}
105
106/// State derives only from canonical snapshots; callers persist each mutation before effects.
107#[derive(Debug, Clone)]
108pub struct MatrixState {
109    snapshot: MatrixSnapshot,
110}
111
112impl MatrixState {
113    /// Create an idle matrix after structural and path validation.
114    pub fn create(spec: MatrixSpec, root: &Path) -> MatrixResult<Self> {
115        validate_spec(&spec)?;
116        validate_workspace(&spec, root)?;
117        let tasks = spec
118            .tasks
119            .iter()
120            .map(|task| MatrixTaskState {
121                id: task.id.clone(),
122                status: MatrixTaskStatus::Queued,
123                attempts: Vec::new(),
124                automatic_retries: 0,
125            })
126            .collect();
127        Ok(Self {
128            snapshot: MatrixSnapshot {
129                spec,
130                lifecycle: MatrixLifecycle::Created,
131                tasks,
132                generation: None,
133                revision: 0,
134            },
135        })
136    }
137
138    /// Validate and restore a retained canonical checkpoint.
139    pub fn from_snapshot(snapshot: MatrixSnapshot) -> MatrixResult<Self> {
140        validate_spec(&snapshot.spec)?;
141        require(
142            snapshot.tasks.len() == snapshot.spec.tasks.len()
143                && snapshot
144                    .tasks
145                    .iter()
146                    .zip(&snapshot.spec.tasks)
147                    .all(|(state, spec)| state.id == spec.id),
148            "snapshot task identities disagree with specification",
149        )?;
150        let mut attempt_ids = BTreeSet::new();
151        let mut worker_ids = BTreeSet::new();
152        let mut active_count = 0;
153        let mut active_writer = false;
154        let mut resource_usage = BTreeMap::<&str, u64>::new();
155        for (task, spec) in snapshot.tasks.iter().zip(&snapshot.spec.tasks) {
156            require(task.automatic_retries <= 1, "snapshot exceeds automatic retry limit")?;
157            let active = task.attempts.iter().filter(|attempt| !attempt.cleanup_confirmed).count();
158            require(active <= 1, "snapshot has overlapping attempts for one task")?;
159            if active == 1 {
160                require(
161                    matches!(task.status, MatrixTaskStatus::Assigned | MatrixTaskStatus::CleanupUncertain),
162                    "snapshot active attempt has inconsistent task status",
163                )?;
164                active_count += 1;
165                active_writer |= spec.access == WorkspaceAccess::Write;
166                for (name, quantity) in &spec.resources {
167                    *resource_usage.entry(name).or_default() += u64::from(*quantity);
168                }
169            } else {
170                require(
171                    !matches!(task.status, MatrixTaskStatus::Assigned | MatrixTaskStatus::CleanupUncertain),
172                    "snapshot assigned task has no owned attempt",
173                )?;
174            }
175            for attempt in &task.attempts {
176                require(
177                    !attempt.id.is_empty()
178                        && !attempt.worker_id.is_empty()
179                        && attempt_ids.insert(&attempt.id)
180                        && worker_ids.insert(&attempt.worker_id),
181                    "snapshot has duplicate or empty attempt identities",
182                )?;
183            }
184            if task.status == MatrixTaskStatus::Verified {
185                let attempt = task
186                    .attempts
187                    .last()
188                    .ok_or_else(|| MatrixError("verified task has no attempt".into()))?;
189                require(
190                    attempt.phase == MatrixPhase::Verify
191                        && attempt.outcome == Some(MatrixOutcome::Success)
192                        && attempt.cleanup_confirmed
193                        && attempt.generation == snapshot.generation,
194                    "verified task has no completed current-generation attempt",
195                )?;
196                validate_evidence(&spec.checks, attempt, &attempt.evidence)?;
197            }
198        }
199        require(
200            active_count <= 5 && (!active_writer || active_count <= 1),
201            "snapshot overcommits worker or workspace leases",
202        )?;
203        require(
204            resource_usage.iter().all(|(name, used)| {
205                snapshot
206                    .spec
207                    .resources
208                    .get(*name)
209                    .is_some_and(|capacity| *used <= u64::from(*capacity))
210            }),
211            "snapshot overcommits named resources",
212        )?;
213        if snapshot.lifecycle == MatrixLifecycle::Succeeded {
214            require(
215                snapshot.tasks.iter().all(|task| {
216                    task.status == MatrixTaskStatus::Verified
217                        && task.attempts.iter().all(|attempt| attempt.cleanup_confirmed)
218                }),
219                "successful matrix has incomplete or unstopped work",
220            )?;
221        }
222        Ok(Self { snapshot })
223    }
224
225    /// Inspect the complete canonical state.
226    pub fn snapshot(&self) -> &MatrixSnapshot {
227        &self.snapshot
228    }
229    /// Build the canonical event for an acknowledged persistence barrier.
230    pub fn event(&self) -> ThreadEvent {
231        ThreadEvent::MatrixUpdated(Box::new(self.snapshot.clone()))
232    }
233    fn changed(&mut self) {
234        self.snapshot.revision = self.snapshot.revision.saturating_add(1);
235    }
236    fn terminal(&self) -> bool {
237        matches!(self.snapshot.lifecycle, MatrixLifecycle::Cancelled | MatrixLifecycle::Succeeded)
238    }
239    /// Assignments whose owned work has not been confirmed stopped.
240    pub fn active_assignments(&self) -> Vec<MatrixAssignment> {
241        self.snapshot
242            .tasks
243            .iter()
244            .flat_map(|task| {
245                task.attempts
246                    .iter()
247                    .filter(|attempt| !attempt.cleanup_confirmed)
248                    .map(|attempt| MatrixAssignment {
249                        task_id: task.id.clone(),
250                        attempt_id: attempt.id.clone(),
251                        worker_id: attempt.worker_id.clone(),
252                        phase: attempt.phase,
253                        generation: attempt.generation.clone(),
254                    })
255            })
256            .collect()
257    }
258    /// Freeze creation and permit scheduler dispatch.
259    pub fn start(&mut self) -> MatrixResult<()> {
260        require(self.snapshot.lifecycle == MatrixLifecycle::Created, "only a created matrix can start")?;
261        self.snapshot.lifecycle = MatrixLifecycle::Running;
262        self.changed();
263        Ok(())
264    }
265    /// Stop new dispatch while admitted workers finish.
266    pub fn pause(&mut self) -> MatrixResult<()> {
267        require(
268            matches!(self.snapshot.lifecycle, MatrixLifecycle::Running | MatrixLifecycle::Verifying),
269            "matrix is not running",
270        )?;
271        self.snapshot.lifecycle = MatrixLifecycle::Paused;
272        self.changed();
273        Ok(())
274    }
275    /// Resume dispatch from a paused phase.
276    pub fn resume(&mut self) -> MatrixResult<()> {
277        require(self.snapshot.lifecycle == MatrixLifecycle::Paused, "matrix is not paused")?;
278        self.snapshot.lifecycle = if self.snapshot.generation.is_some() {
279            MatrixLifecycle::Verifying
280        } else {
281            MatrixLifecycle::Running
282        };
283        self.changed();
284        Ok(())
285    }
286    /// Permanently stop new dispatch; owned workers still require cleanup.
287    pub fn cancel(&mut self) -> MatrixResult<()> {
288        require(!self.terminal(), "matrix is already terminal")?;
289        self.snapshot.lifecycle = MatrixLifecycle::Cancelled;
290        self.changed();
291        Ok(())
292    }
293
294    /// Reserve slot, shared workspace lease, and all named resources as one state mutation.
295    pub fn reserve_ready(&mut self, configured_cap: usize) -> MatrixResult<Vec<MatrixAssignment>> {
296        if !matches!(self.snapshot.lifecycle, MatrixLifecycle::Running | MatrixLifecycle::Verifying) {
297            return Ok(Vec::new());
298        }
299        let cap = configured_cap.min(5);
300        let phase = if self.snapshot.lifecycle == MatrixLifecycle::Verifying {
301            MatrixPhase::Verify
302        } else {
303            MatrixPhase::Execute
304        };
305        let mut assignments = Vec::new();
306        for index in 0..self.snapshot.tasks.len() {
307            let active = self.active_assignments();
308            if active.len() >= cap {
309                break;
310            }
311            let task = self
312                .snapshot
313                .spec
314                .tasks
315                .get(index)
316                .ok_or_else(|| MatrixError("unknown task specification".into()))?;
317            let state = self
318                .snapshot
319                .tasks
320                .get(index)
321                .ok_or_else(|| MatrixError("unknown task index".into()))?;
322            let ready = match phase {
323                MatrixPhase::Execute => {
324                    state.status == MatrixTaskStatus::Queued
325                        && task.dependencies.iter().all(|id| {
326                            self.snapshot.tasks.iter().any(|state| {
327                                &state.id == id
328                                    && matches!(state.status, MatrixTaskStatus::Executed | MatrixTaskStatus::Verified)
329                            })
330                        })
331                }
332                MatrixPhase::Verify => state.status == MatrixTaskStatus::Executed,
333            };
334            if !ready {
335                continue;
336            }
337            let mut usage = BTreeMap::<&str, u64>::new();
338            let mut blocked = false;
339            for assignment in &active {
340                let Some(owner) = self.snapshot.spec.tasks.iter().find(|task| task.id == assignment.task_id) else {
341                    return Err(MatrixError("unknown active task".into()));
342                };
343                if owner.access == WorkspaceAccess::Write || task.access == WorkspaceAccess::Write {
344                    blocked = true;
345                }
346                for (resource, count) in &owner.resources {
347                    let used = usage.entry(resource).or_default();
348                    *used += u64::from(*count);
349                }
350            }
351            if blocked
352                || task.resources.iter().any(|(resource, count)| {
353                    self.snapshot.spec.resources.get(resource).is_none_or(|capacity| {
354                        usage.get(resource.as_str()).copied().unwrap_or(0) + u64::from(*count) > u64::from(*capacity)
355                    })
356                })
357            {
358                continue;
359            }
360            let assignment = MatrixAssignment {
361                task_id: task.id.clone(),
362                attempt_id: uuid::Uuid::new_v4().to_string(),
363                worker_id: uuid::Uuid::new_v4().to_string(),
364                phase,
365                generation: self.snapshot.generation.clone(),
366            };
367            let state = self
368                .snapshot
369                .tasks
370                .get_mut(index)
371                .ok_or_else(|| MatrixError("unknown task index".into()))?;
372            state.status = MatrixTaskStatus::Assigned;
373            state.attempts.push(MatrixAttempt {
374                id: assignment.attempt_id.clone(),
375                worker_id: assignment.worker_id.clone(),
376                phase,
377                generation: assignment.generation.clone(),
378                launch_requested: false,
379                outcome: None,
380                cleanup_confirmed: false,
381                evidence: Vec::new(),
382            });
383            assignments.push(assignment);
384        }
385        if !assignments.is_empty() {
386            self.changed();
387        }
388        Ok(assignments)
389    }
390
391    /// Persist this marker before invoking a process or worker launch.
392    pub fn mark_launch_requested(&mut self, attempt_id: &str) -> MatrixResult<()> {
393        let attempt = self
394            .snapshot
395            .tasks
396            .iter_mut()
397            .flat_map(|task| task.attempts.iter_mut())
398            .find(|attempt| attempt.id == attempt_id)
399            .ok_or_else(|| MatrixError("unknown launch attempt".into()))?;
400        require(attempt.outcome.is_none() && !attempt.launch_requested, "launch was already requested")?;
401        attempt.launch_requested = true;
402        self.changed();
403        Ok(())
404    }
405
406    /// Release a launch reservation after launch failure and confirmed cleanup.
407    pub fn rollback_launch(&mut self, attempt_id: &str) -> MatrixResult<()> {
408        let task = self
409            .snapshot
410            .tasks
411            .iter_mut()
412            .find(|task| task.attempts.last().is_some_and(|attempt| attempt.id == attempt_id))
413            .ok_or_else(|| MatrixError("unknown launch attempt".into()))?;
414        let attempt = task
415            .attempts
416            .last_mut()
417            .ok_or_else(|| MatrixError("missing launch attempt".into()))?;
418        require(attempt.outcome.is_none(), "attempt already reported")?;
419        attempt.cleanup_confirmed = true;
420        attempt.outcome = Some(MatrixOutcome::Interrupted);
421        task.status = if attempt.phase == MatrixPhase::Execute {
422            MatrixTaskStatus::Queued
423        } else {
424            MatrixTaskStatus::Executed
425        };
426        self.changed();
427        Ok(())
428    }
429
430    /// Runtime must resolve evidence IDs to owned completed canonical commands before calling.
431    pub fn report(
432        &mut self,
433        attempt_id: &str,
434        worker_id: &str,
435        outcome: MatrixOutcome,
436        evidence: Vec<MatrixCommandEvidence>,
437        cleanup_confirmed: bool,
438    ) -> MatrixResult<()> {
439        let index = self
440            .snapshot
441            .tasks
442            .iter()
443            .position(|task| task.attempts.last().is_some_and(|attempt| attempt.id == attempt_id))
444            .ok_or_else(|| MatrixError("stale or unknown report attempt".into()))?;
445        let task = self
446            .snapshot
447            .tasks
448            .get(index)
449            .ok_or_else(|| MatrixError("unknown task index".into()))?;
450        let attempt = task
451            .attempts
452            .last()
453            .ok_or_else(|| MatrixError("missing report attempt".into()))?;
454        require(
455            attempt.worker_id == worker_id && attempt.outcome.is_none(),
456            "duplicate report or wrong worker identity",
457        )?;
458        require(task.status == MatrixTaskStatus::Assigned, "task is not assigned")?;
459        if outcome == MatrixOutcome::Success && attempt.phase == MatrixPhase::Verify {
460            let generation = attempt
461                .generation
462                .as_deref()
463                .ok_or_else(|| MatrixError("verification has no generation".into()))?;
464            require(self.snapshot.generation.as_deref() == Some(generation), "stale verification generation")?;
465            let checks = &self
466                .snapshot
467                .spec
468                .tasks
469                .get(index)
470                .ok_or_else(|| MatrixError("unknown task specification".into()))?
471                .checks;
472            validate_evidence(checks, attempt, &evidence)?;
473        }
474        let phase = attempt.phase;
475        let terminal = self.terminal();
476        let task = self
477            .snapshot
478            .tasks
479            .get_mut(index)
480            .ok_or_else(|| MatrixError("unknown task index".into()))?;
481        let attempt = task
482            .attempts
483            .last_mut()
484            .ok_or_else(|| MatrixError("missing report attempt".into()))?;
485        attempt.outcome = Some(outcome);
486        attempt.cleanup_confirmed = cleanup_confirmed;
487        attempt.evidence = evidence;
488        task.status = if !cleanup_confirmed {
489            MatrixTaskStatus::CleanupUncertain
490        } else {
491            match outcome {
492                MatrixOutcome::Success => {
493                    if phase == MatrixPhase::Execute {
494                        MatrixTaskStatus::Executed
495                    } else {
496                        MatrixTaskStatus::Verified
497                    }
498                }
499                MatrixOutcome::Interrupted => MatrixTaskStatus::Interrupted,
500                MatrixOutcome::TimedOut => MatrixTaskStatus::TimedOut,
501                _ => MatrixTaskStatus::Failed,
502            }
503        };
504        if !terminal {
505            if cleanup_confirmed
506                && matches!(outcome, MatrixOutcome::Interrupted | MatrixOutcome::TimedOut)
507                && self
508                    .snapshot
509                    .spec
510                    .tasks
511                    .get(index)
512                    .ok_or_else(|| MatrixError("unknown task specification".into()))?
513                    .replay_safe
514                && task.automatic_retries == 0
515            {
516                task.automatic_retries += 1;
517                task.status = if phase == MatrixPhase::Execute {
518                    MatrixTaskStatus::Queued
519                } else {
520                    MatrixTaskStatus::Executed
521                };
522            } else if outcome != MatrixOutcome::Success || !cleanup_confirmed {
523                self.snapshot.lifecycle = MatrixLifecycle::Blocked;
524            }
525            if self.has_blockers() {
526                self.snapshot.lifecycle = MatrixLifecycle::Blocked;
527            }
528        }
529        self.changed();
530        Ok(())
531    }
532
533    /// Start complete final verification after every execution task finishes.
534    pub fn begin_verification(&mut self, generation: String) -> MatrixResult<()> {
535        require(
536            self.snapshot.lifecycle == MatrixLifecycle::Running
537                && self.active_assignments().is_empty()
538                && self.snapshot.tasks.iter().all(|task| task.status == MatrixTaskStatus::Executed)
539                && !generation.is_empty(),
540            "execution must finish before final verification",
541        )?;
542        self.snapshot.generation = Some(generation);
543        self.snapshot.lifecycle = MatrixLifecycle::Verifying;
544        self.changed();
545        Ok(())
546    }
547    /// The runtime fingerprints again after cleanup before committing success.
548    pub fn finalize_verification(&mut self, generation: &str) -> MatrixResult<()> {
549        require(
550            self.snapshot.lifecycle == MatrixLifecycle::Verifying
551                && self.snapshot.generation.as_deref() == Some(generation)
552                && self.active_assignments().is_empty()
553                && self.snapshot.tasks.iter().all(|task| task.status == MatrixTaskStatus::Verified),
554            "verification is incomplete, paused, or its generation changed",
555        )?;
556        self.snapshot.lifecycle = MatrixLifecycle::Succeeded;
557        self.changed();
558        Ok(())
559    }
560    /// Explicit resume reconciles prepared assignments and blocks ambiguous survivors.
561    pub fn recover(&mut self) -> MatrixResult<()> {
562        let assignments = self.active_assignments();
563        for assignment in assignments {
564            let attempt = self
565                .snapshot
566                .tasks
567                .iter()
568                .flat_map(|task| &task.attempts)
569                .find(|attempt| attempt.id == assignment.attempt_id)
570                .ok_or_else(|| MatrixError("missing recovery attempt".into()))?;
571            if !attempt.launch_requested {
572                self.rollback_launch(&assignment.attempt_id)?;
573            } else if attempt.outcome.is_none() {
574                self.report(
575                    &assignment.attempt_id,
576                    &assignment.worker_id,
577                    MatrixOutcome::Interrupted,
578                    Vec::new(),
579                    false,
580                )?;
581            } else if !self.terminal() {
582                self.snapshot.lifecycle = MatrixLifecycle::Blocked;
583                self.changed();
584            }
585        }
586        Ok(())
587    }
588    /// Invalidate every verification result after a source generation change.
589    pub fn invalidate_verification(&mut self, generation: String) -> MatrixResult<()> {
590        require(
591            !self.terminal()
592                && !self.has_blockers()
593                && self.snapshot.generation.is_some()
594                && self.active_assignments().is_empty()
595                && !generation.is_empty(),
596            "stop owned work before changing verification generation",
597        )?;
598        for task in &mut self.snapshot.tasks {
599            task.status = MatrixTaskStatus::Executed;
600        }
601        self.snapshot.generation = Some(generation);
602        if self.snapshot.lifecycle != MatrixLifecycle::Paused {
603            self.snapshot.lifecycle = MatrixLifecycle::Verifying;
604        }
605        self.changed();
606        Ok(())
607    }
608    /// Retry a failed task after an explicit coordinator decision and cleanup.
609    pub fn retry(&mut self, task_id: &str) -> MatrixResult<()> {
610        require(!self.terminal(), "terminal matrix cannot retry")?;
611        require(self.active_assignments().is_empty(), "wait for owned work to stop before retry")?;
612        let task = self
613            .snapshot
614            .tasks
615            .iter_mut()
616            .find(|task| task.id == task_id)
617            .ok_or_else(|| MatrixError("unknown task".into()))?;
618        require(
619            matches!(
620                task.status,
621                MatrixTaskStatus::Failed | MatrixTaskStatus::Interrupted | MatrixTaskStatus::TimedOut
622            ) && task.attempts.iter().all(|attempt| attempt.cleanup_confirmed),
623            "task cannot retry before confirmed cleanup",
624        )?;
625        task.status = MatrixTaskStatus::Queued;
626        // An explicit retry delegates repairs before a fresh complete check phase.
627        for task in &mut self.snapshot.tasks {
628            if task.status == MatrixTaskStatus::Verified {
629                task.status = MatrixTaskStatus::Executed;
630            }
631        }
632        self.snapshot.generation = None;
633        self.snapshot.lifecycle = MatrixLifecycle::Running;
634        self.changed();
635        Ok(())
636    }
637    /// Cleanup must be established by an ownership token, never a PID alone.
638    pub fn confirm_cleanup(&mut self, attempt_id: &str, worker_id: &str) -> MatrixResult<()> {
639        let index = self
640            .snapshot
641            .tasks
642            .iter()
643            .position(|task| task.attempts.last().is_some_and(|attempt| attempt.id == attempt_id))
644            .ok_or_else(|| MatrixError("unknown cleanup attempt".into()))?;
645        let terminal = self.terminal();
646        let replay_safe = self
647            .snapshot
648            .spec
649            .tasks
650            .get(index)
651            .ok_or_else(|| MatrixError("unknown task specification".into()))?
652            .replay_safe;
653        let task = self
654            .snapshot
655            .tasks
656            .get_mut(index)
657            .ok_or_else(|| MatrixError("unknown task index".into()))?;
658        let attempt = task
659            .attempts
660            .last_mut()
661            .ok_or_else(|| MatrixError("missing cleanup attempt".into()))?;
662        require(attempt.worker_id == worker_id && attempt.outcome.is_some(), "cleanup identity or outcome missing")?;
663        attempt.cleanup_confirmed = true;
664        task.status = match attempt.outcome {
665            Some(MatrixOutcome::Success) => {
666                if attempt.phase == MatrixPhase::Execute {
667                    MatrixTaskStatus::Executed
668                } else {
669                    MatrixTaskStatus::Verified
670                }
671            }
672            Some(MatrixOutcome::Interrupted) => MatrixTaskStatus::Interrupted,
673            Some(MatrixOutcome::TimedOut) => MatrixTaskStatus::TimedOut,
674            _ => MatrixTaskStatus::Failed,
675        };
676        if !terminal
677            && replay_safe
678            && task.automatic_retries == 0
679            && matches!(attempt.outcome, Some(MatrixOutcome::Interrupted | MatrixOutcome::TimedOut))
680        {
681            task.automatic_retries += 1;
682            task.status = if attempt.phase == MatrixPhase::Execute {
683                MatrixTaskStatus::Queued
684            } else {
685                MatrixTaskStatus::Executed
686            };
687            if !self.has_blockers() {
688                self.snapshot.lifecycle = if self.snapshot.generation.is_some() {
689                    MatrixLifecycle::Verifying
690                } else {
691                    MatrixLifecycle::Running
692                };
693            }
694        }
695        self.changed();
696        Ok(())
697    }
698
699    fn has_blockers(&self) -> bool {
700        self.snapshot.tasks.iter().any(|task| {
701            matches!(
702                task.status,
703                MatrixTaskStatus::Failed
704                    | MatrixTaskStatus::Interrupted
705                    | MatrixTaskStatus::TimedOut
706                    | MatrixTaskStatus::CleanupUncertain
707            )
708        })
709    }
710}
711
712fn validate_evidence(
713    checks: &[String],
714    attempt: &MatrixAttempt,
715    evidence: &[MatrixCommandEvidence],
716) -> MatrixResult<()> {
717    let generation = attempt
718        .generation
719        .as_deref()
720        .ok_or_else(|| MatrixError("verification has no generation".into()))?;
721    require(
722        evidence.len() == checks.len()
723            && checks.iter().zip(evidence).all(|(check, record)| {
724                record.command == *check
725                    && record.attempt_id == attempt.id
726                    && record.worker_id == attempt.worker_id
727                    && record.generation == generation
728                    && !record.event_id.is_empty()
729                    && record.exit_code == Some(0)
730                    && !record.cancelled
731            })
732            && evidence.iter().map(|record| &record.event_id).collect::<BTreeSet<_>>().len() == evidence.len(),
733        "missing, unrelated, duplicate, cancelled, failed, or stale verification evidence",
734    )
735}
736
737/// Latest complete durable checkpoint for each matrix, in stable ID order.
738pub fn replay(events: impl IntoIterator<Item = ThreadEvent>) -> MatrixResult<BTreeMap<String, MatrixState>> {
739    let mut matrices = BTreeMap::<String, MatrixState>::new();
740    for event in events {
741        if let ThreadEvent::MatrixUpdated(snapshot) = event {
742            if let Some(previous) = matrices.get(&snapshot.spec.id) {
743                require(previous.snapshot.spec == snapshot.spec, "matrix specification changed after creation")?;
744                require(snapshot.revision > previous.snapshot.revision, "matrix revision did not advance")?;
745            }
746            let state = MatrixState::from_snapshot(*snapshot)?;
747            matrices.insert(state.snapshot.spec.id.clone(), state);
748        }
749    }
750    Ok(matrices)
751}
752
753#[cfg(test)]
754mod tests;