1use 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}")]
9pub struct MatrixError(pub String);
11pub 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
22pub 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
83pub 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#[derive(Debug, Clone)]
108pub struct MatrixState {
109 snapshot: MatrixSnapshot,
110}
111
112impl MatrixState {
113 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 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 pub fn snapshot(&self) -> &MatrixSnapshot {
227 &self.snapshot
228 }
229 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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
737pub 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;