Skip to main content

vtcode_core/subagents/matrix/
mod.rs

1//! Local matrix orchestration beside the subagent runtime. Canonical snapshots
2//! are acknowledged before launch and before exposing completion.
3mod 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/// Access to the existing authoritative drain, not a second persisted store.
24#[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/// Only the runtime can bind this capability to a worker registry.
53#[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    /// Coalesced notification projection; lifecycle authority stays in events.
116    pub fn take_matrix_completion(&self) -> Option<MatrixSnapshot> {
117        self.matrix.completion.lock().take()
118    }
119    /// Idle wakeups inspect readiness without consuming the outer loop's result.
120    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        // Admission stays closed until replay and its reconciliation checkpoint
130        // have been acknowledged. A failed load must never enable discovery.
131        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    /// Headless completion keeps the canonical sink alive until admitted work stops.
170    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        // Stopping owned work cannot depend on a functioning persistence sink.
184        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                // Block new discovery before inspecting its reserved permits,
310                // including when restoring a matrix into a warm session.
311                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                    // Stop this driver's owned work without turning an internal
351                    // failure into the user's terminal cancellation decision.
352                    driver_cancel.cancel();
353                }
354                controller.matrix.driver_active.store(false, Ordering::Release);
355                // Cleanup-uncertain state keeps the coordinator restricted.
356                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                // Preserve stopped-work results even if inputs became unavailable.
397                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;