vtcode_core/subagents/matrix/
mod.rs1mod fingerprint;
4mod projection;
5mod worker;
6
7use super::SubagentController;
8use crate::exec::events::matrix::*;
9use anyhow::{Context, Result, bail, ensure};
10use futures::future::BoxFuture;
11use parking_lot::RwLock;
12use serde_json::{Value, json};
13use std::collections::BTreeMap;
14use std::sync::{
15 Arc,
16 atomic::{AtomicBool, Ordering},
17};
18use tokio::sync::{Mutex, Notify};
19use tokio::task::JoinSet;
20use tokio_util::sync::CancellationToken;
21use vtcode_memory::matrix::MatrixState;
22
23#[derive(Clone)]
25pub struct MatrixPersistence {
26 pub persist: Arc<dyn Fn(MatrixSnapshot) -> BoxFuture<'static, Result<()>> + Send + Sync>,
27 pub load: Arc<dyn Fn() -> BoxFuture<'static, Result<Vec<MatrixSnapshot>>> + Send + Sync>,
28}
29
30#[derive(Default)]
31pub(super) struct MatrixRuntime {
32 #[cfg(test)]
33 executor_override: RwLock<Option<TestExecutor>>,
34 persistence: RwLock<Option<MatrixPersistence>>,
35 updated_at: RwLock<chrono::DateTime<chrono::Utc>>,
36 state: Mutex<Option<MatrixState>>,
37 driver_active: AtomicBool,
38 pub(super) executing: AtomicBool,
39 notify: Notify,
40 pub(super) cancellation: RwLock<CancellationToken>,
41 error: RwLock<Option<String>>,
42 completion: parking_lot::Mutex<Option<MatrixSnapshot>>,
43}
44
45#[cfg(test)]
46type TestExecutor = Arc<
47 dyn Fn(MatrixAssignment, MatrixTaskSpec, CancellationToken) -> BoxFuture<'static, worker::WorkerResult>
48 + Send
49 + Sync,
50>;
51
52#[derive(Clone)]
54pub(crate) struct MatrixWorkerContext {
55 pub assignment: MatrixAssignment,
56 report: Arc<parking_lot::Mutex<Option<Value>>>,
57}
58impl MatrixWorkerContext {
59 fn new(assignment: MatrixAssignment) -> Self {
60 Self {
61 assignment,
62 report: Arc::new(parking_lot::Mutex::new(None)),
63 }
64 }
65 pub(crate) fn reported_outcome(&self) -> Option<MatrixOutcome> {
66 self.report
67 .lock()
68 .as_ref()
69 .and_then(|report| report.get("outcome").and_then(Value::as_str))
70 .and_then(|outcome| match outcome {
71 "executed" => Some(MatrixOutcome::Success),
72 "failed" => Some(MatrixOutcome::Failed),
73 "permission_denied" => Some(MatrixOutcome::PermissionDenied),
74 "budget_exhausted" => Some(MatrixOutcome::BudgetExhausted),
75 "interrupted" => Some(MatrixOutcome::Interrupted),
76 "timed_out" => Some(MatrixOutcome::TimedOut),
77 _ => None,
78 })
79 }
80 pub(crate) fn report(&self, args: Value) -> Result<Value> {
81 ensure!(args.get("action").and_then(Value::as_str) == Some("report"), "matrix workers may only report");
82 ensure!(
83 args.get("matrix_id").is_none()
84 && args.get("task_id").is_none()
85 && args.get("attempt_id").is_none()
86 && args.get("worker_id").is_none(),
87 "worker report identity is runtime-owned"
88 );
89 let outcome = args
90 .get("outcome")
91 .and_then(Value::as_str)
92 .context("matrix report requires outcome")?;
93 ensure!(
94 [
95 "executed",
96 "failed",
97 "interrupted",
98 "timed_out",
99 "permission_denied",
100 "budget_exhausted"
101 ]
102 .contains(&outcome),
103 "invalid matrix report outcome"
104 );
105 let mut report = self.report.lock();
106 ensure!(report.is_none(), "duplicate matrix report");
107 *report = Some(args);
108 Ok(
109 json!({"accepted":true,"task_id":self.assignment.task_id,"attempt_id":self.assignment.attempt_id,"final_success":false}),
110 )
111 }
112}
113
114impl SubagentController {
115 pub fn take_matrix_completion(&self) -> Option<MatrixSnapshot> {
117 self.matrix.completion.lock().take()
118 }
119 pub fn has_matrix_completion(&self) -> bool {
121 self.matrix.completion.lock().is_some()
122 }
123 pub async fn set_matrix_persistence(&self, persistence: MatrixPersistence) -> Result<()> {
124 let mut guard = self.matrix.state.lock().await;
125 *self.matrix.persistence.write() = Some(persistence.clone());
126 if guard.is_some() {
127 return Ok(());
128 }
129 self.matrix.executing.store(true, Ordering::Release);
132 let mut retained = (persistence.load)().await?.into_iter().filter(|snapshot| {
133 !matches!(snapshot.lifecycle, MatrixLifecycle::Cancelled | MatrixLifecycle::Succeeded)
134 || snapshot
135 .tasks
136 .iter()
137 .flat_map(|task| &task.attempts)
138 .any(|attempt| !attempt.cleanup_confirmed)
139 });
140 let snapshot = retained.next();
141 ensure!(retained.next().is_none(), "multiple unreconciled matrices in one local session");
142 if let Some(snapshot) = snapshot {
143 let mut restored = MatrixState::from_snapshot(snapshot)?;
144 self.matrix.recover(&mut restored).await?;
145 self.matrix.update_execution_admission(&restored);
146 *guard = Some(restored);
147 } else {
148 self.matrix.executing.store(false, Ordering::Release);
149 }
150 Ok(())
151 }
152 pub fn matrix_is_executing(&self) -> bool {
153 self.matrix.executing.load(Ordering::Acquire)
154 }
155 pub(super) fn ensure_ordinary_delegation_allowed(&self) -> Result<()> {
156 ensure!(!self.matrix_is_executing(), "active matrix execution must use scheduler-owned workers");
157 Ok(())
158 }
159 pub(crate) fn matrix_is_driving(&self) -> bool {
160 self.matrix.driver_active.load(Ordering::Acquire)
161 }
162 pub(crate) fn request_matrix_stop(&self) {
163 self.matrix.cancellation.read().cancel();
164 self.matrix.notify.notify_one();
165 }
166 pub(crate) async fn matrix_snapshot(&self) -> Option<MatrixSnapshot> {
167 self.matrix.state.lock().await.as_ref().map(|state| state.snapshot().clone())
168 }
169 pub(crate) async fn wait_matrix_idle(&self) -> Option<MatrixSnapshot> {
171 loop {
172 let notified = self.background_completion_notify.notified();
173 if !self.matrix.driver_active.load(Ordering::Acquire) {
174 return self.matrix_snapshot().await;
175 }
176 tokio::select! {
177 () = notified => {}
178 () = tokio::time::sleep(std::time::Duration::from_millis(250)) => {}
179 }
180 }
181 }
182 pub async fn cancel_matrix(&self) -> Result<()> {
183 self.request_matrix_stop();
185 let mut guard = self.matrix.state.lock().await;
186 if let Some(state) = guard.as_ref()
187 && !matches!(state.snapshot().lifecycle, MatrixLifecycle::Cancelled | MatrixLifecycle::Succeeded)
188 {
189 let mut next = state.clone();
190 next.cancel()?;
191 self.matrix.persist(&next).await?;
192 self.matrix.update_execution_admission(&next);
193 *guard = Some(next);
194 }
195 Ok(())
196 }
197 pub async fn matrix_control(&self, args: Value) -> Result<Value> {
198 let action = args.get("action").and_then(Value::as_str).context("matrix requires action")?;
199 ensure!(action != "report", "matrix report is worker-only");
200 if ["start", "resume", "retry"].contains(&action) {
201 ensure!(
202 self.config.depth == 0
203 && self.config.vt_cfg.subagents.enabled
204 && self.config.vt_cfg.subagents.max_concurrent > 0
205 && self.config.vt_cfg.subagents.max_depth > 0,
206 "matrix requires root-session enabled subagents with positive capacity and depth"
207 );
208 ensure!(!self.matrix.cancellation.read().is_cancelled(), "user cancellation prevents matrix continuation");
209 }
210 let requested_id = args.get("matrix_id").and_then(Value::as_str);
211 let mut guard = self.matrix.state.lock().await;
212 if guard.is_none() && action != "create" {
213 let persistence = self
214 .matrix
215 .persistence
216 .read()
217 .clone()
218 .context("canonical matrix persistence unavailable")?;
219 let snapshots = (persistence.load)().await?;
220 let snapshot = snapshots
221 .into_iter()
222 .find(|snapshot| Some(snapshot.spec.id.as_str()) == requested_id)
223 .context("unknown matrix; provide matrix_id")?;
224 let mut restored = MatrixState::from_snapshot(snapshot)?;
225 self.matrix.update_execution_admission(&restored);
226 self.matrix.recover(&mut restored).await?;
227 *guard = Some(restored);
228 }
229 if action == "create" {
230 ensure!(!self.matrix.driver_active.load(Ordering::Acquire), "matrix driver is still cleaning up");
231 if let Some(state) = guard.as_ref() {
232 ensure!(
233 matches!(state.snapshot().lifecycle, MatrixLifecycle::Cancelled | MatrixLifecycle::Succeeded)
234 && state.active_assignments().is_empty(),
235 "one local matrix may be active at a time"
236 );
237 }
238 let spec: MatrixSpec =
239 serde_json::from_value(args.get("spec").cloned().context("matrix create requires spec")?)?;
240 ensure!(requested_id.is_none_or(|id| id == spec.id), "matrix_id does not match specification");
241 let persistence = self
242 .matrix
243 .persistence
244 .read()
245 .clone()
246 .context("canonical matrix persistence unavailable")?;
247 let existing = (persistence.load)().await?;
248 ensure!(
249 existing.iter().all(|snapshot| matches!(
250 snapshot.lifecycle,
251 MatrixLifecycle::Cancelled | MatrixLifecycle::Succeeded
252 ) && snapshot
253 .tasks
254 .iter()
255 .flat_map(|task| &task.attempts)
256 .all(|attempt| attempt.cleanup_confirmed)),
257 "resume the existing matrix and reconcile owned cleanup before creating another"
258 );
259 ensure!(
260 !existing.iter().any(|snapshot| snapshot.spec.id == spec.id),
261 "matrix ID already exists; use a new ID"
262 );
263 let state = MatrixState::create(spec, &self.config.workspace_root)?;
264 self.matrix.persist(&state).await?;
265 self.matrix.update_execution_admission(&state);
266 *guard = Some(state);
267 *self.matrix.cancellation.write() = CancellationToken::new();
268 *self.matrix.error.write() = None;
269 } else if action != "status" {
270 let mut state = guard.as_ref().context("matrix is unavailable")?.clone();
271 ensure!(requested_id == Some(state.snapshot().spec.id.as_str()), "matrix_id does not match active matrix");
272 if ["start", "resume", "retry"].contains(&action) && !self.matrix_is_driving() {
273 self.matrix.recover(&mut state).await?;
274 self.matrix.update_execution_admission(&state);
275 *guard = Some(state.clone());
276 }
277 let mut next = state.clone();
278 match action {
279 "start" => {
280 ensure!(
281 self.config.vt_cfg.subagents.enabled && self.config.vt_cfg.subagents.max_concurrent > 0,
282 "matrix requires enabled subagents with positive concurrency"
283 );
284 vtcode_memory::matrix::validate_workspace(&next.snapshot().spec, &self.config.workspace_root)?;
285 fingerprint::fingerprint(&self.config.workspace_root, &next.snapshot().spec).await?;
286 next.start()?;
287 }
288 "pause" => next.pause()?,
289 "resume" => {
290 if next.snapshot().lifecycle == MatrixLifecycle::Paused {
291 next.resume()?;
292 } else {
293 ensure!(
294 matches!(next.snapshot().lifecycle, MatrixLifecycle::Running | MatrixLifecycle::Verifying),
295 "matrix needs a coordinator retry or confirmed owned cleanup"
296 );
297 }
298 }
299 "retry" => {
300 next.retry(args.get("task_id").and_then(Value::as_str).context("retry requires task_id")?)?
301 }
302 "cancel" => {
303 self.matrix.cancellation.read().cancel();
304 next.cancel()?;
305 }
306 _ => bail!("unknown matrix action {action}"),
307 }
308 if ["start", "resume", "retry"].contains(&action) && !self.matrix.driver_active.load(Ordering::Acquire) {
309 self.matrix.executing.store(true, Ordering::Release);
312 if self.admission.available_permits()
313 != self
314 .config
315 .vt_cfg
316 .subagents
317 .max_concurrent
318 .min(vtcode_config::subagents::SUBAGENT_HARD_CONCURRENCY_LIMIT)
319 {
320 self.matrix.update_execution_admission(&state);
321 bail!("wait for discovery workers to stop before matrix dispatch");
322 }
323 }
324 if next.snapshot() != state.snapshot()
325 && let Err(error) = self.matrix.persist(&next).await
326 {
327 self.matrix.update_execution_admission(&state);
328 return Err(error);
329 }
330 self.matrix.update_execution_admission(&next);
331 *guard = Some(next);
332 if action == "cancel" {
333 self.matrix.cancellation.read().cancel();
334 }
335 }
336 let snapshot = guard.as_ref().context("matrix unavailable")?.snapshot().clone();
337 if let Some(id) = requested_id {
338 ensure!(id == snapshot.spec.id, "matrix_id does not match active matrix");
339 }
340 drop(guard);
341 self.matrix.notify.notify_one();
342 if ["start", "resume", "retry"].contains(&action) && !self.matrix.driver_active.swap(true, Ordering::AcqRel) {
343 self.matrix.executing.store(true, Ordering::Release);
344 *self.matrix.error.write() = None;
345 let driver_cancel = self.matrix.cancellation.read().child_token();
346 let controller = self.clone();
347 tokio::spawn(async move {
348 if let Err(error) = controller.matrix_drive(&driver_cancel).await {
349 *controller.matrix.error.write() = Some(format!("{error:#}"));
350 driver_cancel.cancel();
353 }
354 controller.matrix.driver_active.store(false, Ordering::Release);
355 let state = controller.matrix.state.lock().await;
357 if let Some(state) = state.as_ref()
358 && matches!(state.snapshot().lifecycle, MatrixLifecycle::Succeeded | MatrixLifecycle::Blocked)
359 {
360 *controller.matrix.completion.lock() = Some(state.snapshot().clone());
361 }
362 if let Some(state) = state.as_ref() {
363 controller.matrix.update_execution_admission(state);
364 }
365 controller.background_completion_notify.notify_one();
366 });
367 }
368 let mut result = projection::tracker_result(&snapshot);
369 result["matrix"] = json!(snapshot);
370 result["persistence_error"] = json!(self.matrix.error.read().clone());
371 Ok(result)
372 }
373
374 async fn matrix_drive(&self, driver_cancel: &CancellationToken) -> Result<()> {
375 let mut workers = JoinSet::new();
376 let mut pending = BTreeMap::<String, worker::WorkerResult>::new();
377 loop {
378 let mut guard = self.matrix.state.lock().await;
379 let state = guard.as_ref().context("matrix disappeared")?;
380 let mut next = state.clone();
381 if self.matrix.cancellation.read().is_cancelled()
382 && !matches!(next.snapshot().lifecycle, MatrixLifecycle::Cancelled | MatrixLifecycle::Succeeded)
383 {
384 next.cancel()?;
385 }
386 if workers.is_empty() && !pending.is_empty() {
387 for (_, result) in std::mem::take(&mut pending) {
388 next.report(
389 &result.assignment.attempt_id,
390 &result.assignment.worker_id,
391 result.outcome,
392 result.evidence,
393 result.cleanup_confirmed,
394 )?;
395 }
396 self.matrix.persist(&next).await?;
398 *guard = Some(next.clone());
399 if next.snapshot().lifecycle != MatrixLifecycle::Cancelled {
400 let current = fingerprint::fingerprint(&self.config.workspace_root, &next.snapshot().spec).await?;
401 let generation_changed = next
402 .snapshot()
403 .generation
404 .as_deref()
405 .is_some_and(|generation| generation != current);
406 if generation_changed
407 && next.active_assignments().is_empty()
408 && matches!(next.snapshot().lifecycle, MatrixLifecycle::Verifying | MatrixLifecycle::Paused)
409 {
410 next.invalidate_verification(current)?;
411 }
412 }
413 }
414 if workers.is_empty()
415 && next.snapshot().lifecycle == MatrixLifecycle::Running
416 && next
417 .snapshot()
418 .tasks
419 .iter()
420 .all(|task| task.status == MatrixTaskStatus::Executed)
421 {
422 let generation = fingerprint::fingerprint(&self.config.workspace_root, &next.snapshot().spec).await?;
423 next.begin_verification(generation)?;
424 }
425 if workers.is_empty()
426 && next.snapshot().lifecycle == MatrixLifecycle::Verifying
427 && next
428 .snapshot()
429 .tasks
430 .iter()
431 .all(|task| task.status == MatrixTaskStatus::Verified)
432 {
433 let generation = fingerprint::fingerprint(&self.config.workspace_root, &next.snapshot().spec).await?;
434 if next.snapshot().generation.as_deref() != Some(generation.as_str()) {
435 next.invalidate_verification(generation)?;
436 } else {
437 next.finalize_verification(&generation)?;
438 }
439 }
440 let available = self.admission.available_permits();
441 let active = next.active_assignments().len();
442 if matches!(next.snapshot().lifecycle, MatrixLifecycle::Running | MatrixLifecycle::Verifying) {
443 vtcode_memory::matrix::validate_workspace(&next.snapshot().spec, &self.config.workspace_root)?;
444 }
445 let assignments =
446 next.reserve_ready((active + available).min(self.config.vt_cfg.subagents.max_concurrent).min(5))?;
447 let mut launches = Vec::new();
448 for assignment in assignments {
449 match Arc::clone(&self.admission).try_acquire_owned() {
450 Ok(permit) => {
451 let task = next
452 .snapshot()
453 .spec
454 .tasks
455 .iter()
456 .find(|task| task.id == assignment.task_id)
457 .context("unknown assigned task")?
458 .clone();
459 launches.push((assignment, task, permit));
460 }
461 Err(_) => next.rollback_launch(&assignment.attempt_id)?,
462 }
463 }
464 if next.snapshot() != guard.as_ref().context("matrix unavailable")?.snapshot() {
465 self.matrix.persist(&next).await?;
466 *guard = Some(next);
467 self.background_completion_notify.notify_one();
468 }
469 if !launches.is_empty() {
470 let mut launching = guard.as_ref().context("matrix unavailable")?.clone();
471 for (assignment, _, _) in &launches {
472 launching.mark_launch_requested(&assignment.attempt_id)?;
473 }
474 self.matrix.persist(&launching).await?;
475 *guard = Some(launching);
476 }
477 let lifecycle = guard.as_ref().context("matrix unavailable")?.snapshot().lifecycle;
478 drop(guard);
479 for (assignment, task, permit) in launches {
480 let controller = self.clone();
481 let cancel = driver_cancel.child_token();
482 #[cfg(test)]
483 let executor = self.matrix.executor_override.read().clone();
484 workers.spawn(async move {
485 let _permit = permit;
486 #[cfg(test)]
487 if let Some(executor) = executor {
488 return executor(assignment, task, cancel).await;
489 }
490 Box::pin(worker::execute(&controller, assignment, task, cancel)).await
491 });
492 }
493 if workers.is_empty() && !matches!(lifecycle, MatrixLifecycle::Running | MatrixLifecycle::Verifying) {
494 return Ok(());
495 }
496 if workers.is_empty() && lifecycle == MatrixLifecycle::Verifying {
497 continue;
498 }
499 tokio::select! {
500 result = workers.join_next(), if !workers.is_empty() => {
501 let result = result.context("matrix worker missing")?.context("matrix worker panicked; cleanup ownership uncertain")?;
502 if result.assignment.phase == MatrixPhase::Verify {
503 pending.insert(result.assignment.attempt_id.clone(), result);
504 } else {
505 let mut guard = self.matrix.state.lock().await;
506 let mut next = guard.as_ref().context("matrix unavailable")?.clone();
507 next.report(&result.assignment.attempt_id, &result.assignment.worker_id, result.outcome, result.evidence, result.cleanup_confirmed)?;
508 self.matrix.persist(&next).await?;
509 *guard = Some(next);
510 self.background_completion_notify.notify_one();
511 }
512 }
513 () = self.matrix.notify.notified() => {}
514 () = tokio::time::sleep(std::time::Duration::from_millis(250)) => {}
515 }
516 }
517 }
518}
519impl MatrixRuntime {
520 async fn recover(&self, state: &mut MatrixState) -> Result<()> {
521 let previous = state.snapshot().clone();
522 state.recover()?;
523 if state.snapshot() != &previous {
524 self.persist(state).await?;
525 }
526 Ok(())
527 }
528
529 fn update_execution_admission(&self, state: &MatrixState) {
530 let held = !state.active_assignments().is_empty()
531 || !matches!(
532 state.snapshot().lifecycle,
533 MatrixLifecycle::Created | MatrixLifecycle::Succeeded | MatrixLifecycle::Cancelled
534 );
535 self.executing.store(held, Ordering::Release);
536 }
537
538 async fn persist(&self, state: &MatrixState) -> Result<()> {
539 let persistence = self
540 .persistence
541 .read()
542 .clone()
543 .context("canonical matrix persistence unavailable")?;
544 (persistence.persist)(state.snapshot().clone()).await?;
545 *self.updated_at.write() = chrono::Utc::now();
546 Ok(())
547 }
548}
549
550#[cfg(test)]
551mod tests;