Skip to main content

vtcode_core/subagents/
controller_background_ops.rs

1#![allow(
2    unused_imports,
3    reason = "Intentional compatibility, platform, or test-only suppression."
4)]
5use anyhow::{Context, Result, anyhow, bail};
6use chrono::Utc;
7use futures::future::select_all;
8use std::collections::VecDeque;
9use std::path::PathBuf;
10use std::sync::Arc;
11use std::sync::atomic::{AtomicBool, Ordering};
12use tokio::sync::{Notify, RwLock};
13
14use crate::config::VTCodeConfig;
15use crate::config::types::ReasoningEffortLevel;
16use crate::core::agent::runner::{AgentRunner, RunnerSettings};
17use crate::core::agent::task::Task;
18use crate::core::threads::{ThreadBootstrap, ThreadId, ThreadRuntimeHandle, ThreadSnapshot};
19use crate::hooks::{LifecycleHookEngine, SessionStartTrigger};
20use crate::llm::provider::Message;
21use crate::tools::exec_session::{ExecSessionCompletionEvent, ExecSessionManager};
22use crate::tools::pty::{PtyManager, PtySize};
23use crate::utils::session_archive::{SessionArchive, find_session_by_identifier};
24use vtcode_config::SubagentSpec;
25use vtcode_config::auth::OpenAIChatGptAuthHandle;
26
27use self::background::*;
28use self::config::*;
29use self::constants::*;
30use self::discovery::discover_controller_subagents;
31use self::model::*;
32use vtcode_config::subagents::SUBAGENT_HARD_CONCURRENCY_LIMIT;
33
34#[allow(
35    unused_imports,
36    reason = "Intentional compatibility, platform, or test-only suppression."
37)]
38use super::*;
39
40const BACKGROUND_COMPLETION_IDENTITY_CAPACITY: usize = 256;
41
42impl SubagentController {
43    fn clone_for_background_completion_monitor(&self) -> Self {
44        Self {
45            admission: Arc::clone(&self.admission),
46            matrix: Arc::clone(&self.matrix),
47            config: Arc::clone(&self.config),
48            parent_session_id: Arc::clone(&self.parent_session_id),
49            lifecycle_hooks: self.lifecycle_hooks.clone(),
50            state: Arc::clone(&self.state),
51            shutdown_requested: Arc::clone(&self.shutdown_requested),
52            closing: Arc::clone(&self.closing),
53            background_completion_channel: Arc::clone(&self.background_completion_channel),
54            background_completion_notify: Arc::clone(&self.background_completion_notify),
55            background_completion_shutdown: self.background_completion_shutdown.clone(),
56            background_completion_monitor: Arc::clone(&self.background_completion_monitor),
57            background_completion_owners: Arc::clone(&self.background_completion_owners),
58            background_completion_monitor_owner: false,
59        }
60    }
61
62    pub(super) async fn start_background_completion_monitor(&self) {
63        let mut completion_rx = self.config.exec_sessions.subscribe_completion();
64        let controller = self.clone_for_background_completion_monitor();
65        let shutdown = self.background_completion_shutdown.clone();
66        let monitor = tokio::spawn(async move {
67            loop {
68                tokio::select! {
69                    biased;
70                    _ = shutdown.cancelled() => break,
71                    result = completion_rx.recv() => match result {
72                        Ok(event) => {
73                            if let Err(error) = controller.handle_exec_session_completion(event).await {
74                                tracing::warn!(error = %error, "Background completion handling failed");
75                            }
76                        }
77                    Err(broadcast::error::RecvError::Lagged(skipped)) => {
78                            tracing::warn!(skipped, "Background completion monitor lagged; reconciling records");
79                            if let Err(error) = controller.refresh_background_processes().await {
80                                tracing::warn!(error = %error, "Background completion reconciliation failed");
81                            }
82                        }
83                        Err(broadcast::error::RecvError::Closed) => break,
84                    },
85                }
86            }
87        });
88        let mut monitor_slot = self.background_completion_monitor.lock().await;
89        if let Some(previous) = monitor_slot.replace(monitor) {
90            previous.abort();
91            let _ = previous.await;
92        }
93    }
94
95    pub(super) async fn stop_background_completion_monitor(&self) {
96        self.background_completion_shutdown.cancel();
97        let monitor = self.background_completion_monitor.lock().await.take();
98        if let Some(monitor) = monitor {
99            monitor.abort();
100            let _ = monitor.await;
101        }
102    }
103
104    async fn handle_exec_session_completion(&self, event: ExecSessionCompletionEvent) -> Result<()> {
105        if self.shutdown_requested.load(Ordering::Relaxed) {
106            return Ok(());
107        }
108        if !event.managed_background {
109            return Ok(());
110        }
111
112        let record_id = {
113            let state = self.state.read().await;
114            let Some(record_id) = state
115                .background_children
116                .values()
117                .find(|record| record.exec_session_id == event.session_id.as_str())
118                .map(|record| record.id.clone())
119            else {
120                // The session may be a user-managed background command or a
121                // stale completion from a previous controller incarnation.
122                return Ok(());
123            };
124            record_id
125        };
126
127        let Some(snapshot) = self.config.exec_sessions.snapshot_session(event.session_id.as_str()).await.ok() else {
128            return Ok(());
129        };
130        let respawn = self.update_background_record_state(&record_id, Some(snapshot)).await?;
131        if let Some((agent_name, stable_id, restart_attempts)) = respawn {
132            self.ensure_background_record_running(
133                agent_name.as_str(),
134                Some(stable_id.as_str()),
135                restart_attempts,
136                None,
137            )
138            .await?;
139        }
140        self.refresh_background_archive_metadata(&record_id).await?;
141        self.save_background_state().await?;
142
143        self.publish_terminal_background_completion(&record_id, event.session_id.as_str(), Some(event.exit_code))
144            .await
145    }
146
147    async fn publish_terminal_background_completion(
148        &self,
149        record_id: &str,
150        expected_exec_session_id: &str,
151        exit_code: Option<i32>,
152    ) -> Result<()> {
153        if self.shutdown_requested.load(Ordering::Relaxed) {
154            return Ok(());
155        }
156
157        let completion = {
158            let mut state = self.state.write().await;
159            let (task_id, status, summary, error, session_id, exec_session_id, archive_path, transcript_path) = {
160                let Some(record) = state.background_children.get(record_id) else {
161                    return Ok(());
162                };
163                if record.exec_session_id != expected_exec_session_id
164                    || !matches!(record.status, BackgroundSubprocessStatus::Stopped | BackgroundSubprocessStatus::Error)
165                {
166                    return Ok(());
167                }
168                (
169                    record.id.clone(),
170                    record.status,
171                    record.summary.clone(),
172                    record.error.clone(),
173                    record.session_id.clone(),
174                    record.exec_session_id.clone(),
175                    record.archive_path.clone(),
176                    record.transcript_path.clone(),
177                )
178            };
179
180            let identity = format!("{task_id}:{expected_exec_session_id}");
181            if state.background_completion_identities.iter().any(|seen| seen == &identity) {
182                return Ok(());
183            }
184            if state.background_completion_identities.len() >= BACKGROUND_COMPLETION_IDENTITY_CAPACITY {
185                state.background_completion_identities.pop_front();
186            }
187            state.background_completion_identities.push_back(identity);
188
189            BackgroundCompletionEvent {
190                task_id,
191                status,
192                summary,
193                error,
194                session_id,
195                exec_session_id,
196                archive_path,
197                transcript_path,
198                termination_requested: state
199                    .background_children
200                    .get(record_id)
201                    .is_some_and(|record| record.termination_requested),
202                exit_code,
203            }
204        };
205
206        self.background_completion_channel.lock().publish(completion);
207        self.background_completion_notify.notify_one();
208        Ok(())
209    }
210
211    /// Returns status entries for all tracked background subprocesses.
212    pub async fn background_status_entries(&self) -> Vec<BackgroundSubprocessEntry> {
213        let state = self.state.read().await;
214        state
215            .background_children
216            .values()
217            .map(BackgroundRecord::build_status_entry)
218            .collect()
219    }
220
221    /// Returns a snapshot of a background subprocess including its preview output.
222    pub async fn background_snapshot(&self, target: &str) -> Result<BackgroundSubprocessSnapshot> {
223        let _ = self.refresh_background_processes().await?;
224
225        let entry = {
226            let state = self.state.read().await;
227            state
228                .background_children
229                .get(target)
230                .ok_or_else(|| anyhow!("Unknown background subprocess {target}"))?
231                .build_status_entry()
232        };
233
234        let preview = if entry.exec_session_id.is_empty() {
235            String::new()
236        } else {
237            match self
238                .config
239                .exec_sessions
240                .read_session_output(&entry.exec_session_id, false)
241                .await
242            {
243                Ok(Some(output)) => extract_tail_lines(&output, SUBAGENT_PREVIEW_LINES),
244                Ok(None) | Err(_) => {
245                    if let Some(path) = entry.transcript_path.as_ref().or(entry.archive_path.as_ref()) {
246                        load_archive_preview(path).await.unwrap_or_default()
247                    } else {
248                        String::new()
249                    }
250                }
251            }
252        };
253
254        Ok(BackgroundSubprocessSnapshot { entry, preview })
255    }
256
257    /// Returns whether background subagents are enabled in the configuration.
258    #[must_use]
259    pub fn background_subagents_enabled(&self) -> bool {
260        self.config.vt_cfg.subagents.background.enabled
261    }
262
263    /// Returns the configured default background agent name, if any.
264    #[must_use]
265    pub fn configured_default_background_agent(&self) -> Option<&str> {
266        self.config
267            .vt_cfg
268            .subagents
269            .background
270            .default_agent
271            .as_deref()
272            .map(str::trim)
273            .filter(|agent| !agent.is_empty())
274    }
275
276    /// Toggles the default background subagent between running and stopped.
277    pub async fn toggle_default_background_subagent(&self) -> Result<BackgroundSubprocessEntry> {
278        if !self.background_subagents_enabled() {
279            bail!("Background subagents are disabled by configuration");
280        }
281
282        let agent_name = self
283            .configured_default_background_agent()
284            .ok_or_else(|| anyhow!("No default background subagent is configured"))?
285            .to_string();
286        let target_id = background_record_id(agent_name.as_str());
287        let should_stop = {
288            let state = self.state.read().await;
289            state
290                .background_children
291                .get(&target_id)
292                .is_some_and(|record| record.desired_enabled && record.status.is_active())
293        };
294
295        if should_stop {
296            self.graceful_stop_background(&target_id).await
297        } else {
298            self.ensure_background_record_running(agent_name.as_str(), Some(target_id.as_str()), 0, None)
299                .await
300        }
301    }
302
303    /// Restarts background subagents that were previously enabled but are no longer running.
304    pub async fn restore_background_subagents(&self) -> Result<Vec<BackgroundSubprocessEntry>> {
305        let desired_records = {
306            let state = self.state.read().await;
307            state
308                .background_children
309                .values()
310                .filter(|record| record.desired_enabled)
311                .map(|record| {
312                    (
313                        record.id.clone(),
314                        record.agent_name.clone(),
315                        record.exec_session_id.clone(),
316                        record.restart_attempts,
317                    )
318                })
319                .collect::<Vec<_>>()
320        };
321
322        for (record_id, agent_name, exec_session_id, restart_attempts) in desired_records {
323            let is_live = !exec_session_id.is_empty()
324                && self
325                    .config
326                    .exec_sessions
327                    .snapshot_session(&exec_session_id)
328                    .await
329                    .ok()
330                    .is_some_and(|snapshot| exec_session_is_running(&snapshot));
331
332            if is_live || !self.config.vt_cfg.subagents.background.auto_restore {
333                continue;
334            }
335            tracing::info!(
336                agent_name = agent_name.as_str(),
337                record_id = record_id.as_str(),
338                "Restoring background subagent subprocess"
339            );
340            self.ensure_background_record_running(
341                agent_name.as_str(),
342                Some(record_id.as_str()),
343                restart_attempts,
344                None,
345            )
346            .await?;
347        }
348
349        self.refresh_background_processes().await
350    }
351
352    /// Refreshes the state of all background subprocesses and respawns as needed.
353    pub async fn refresh_background_processes(&self) -> Result<Vec<BackgroundSubprocessEntry>> {
354        let record_ids = {
355            let state = self.state.read().await;
356            state.background_children.keys().cloned().collect::<Vec<_>>()
357        };
358
359        let mut changed = false;
360        for record_id in record_ids {
361            let (snapshot_target, before_status, before_error, before_summary, before_desired_enabled) = {
362                let state = self.state.read().await;
363                let record = state.background_children.get(&record_id);
364                (
365                    record.map(|r| r.exec_session_id.clone()),
366                    record.map(|r| r.status),
367                    record.and_then(|r| r.error.clone()),
368                    record.and_then(|r| r.summary.clone()),
369                    record.is_some_and(|r| r.desired_enabled),
370                )
371            };
372
373            let snapshot = if let Some(exec_session_id) = snapshot_target.as_ref()
374                && !exec_session_id.is_empty()
375            {
376                self.config.exec_sessions.snapshot_session(exec_session_id).await.ok()
377            } else {
378                None
379            };
380
381            let respawn = self.update_background_record_state(&record_id, snapshot).await?;
382
383            if let Some((agent_name, stable_id, restart_attempts)) = respawn {
384                self.ensure_background_record_running(
385                    agent_name.as_str(),
386                    Some(stable_id.as_str()),
387                    restart_attempts,
388                    None,
389                )
390                .await?;
391            }
392
393            let changed_this_record = {
394                let state = self.state.read().await;
395                state.background_children.get(&record_id).is_some_and(|r| {
396                    r.status != before_status.unwrap_or(BackgroundSubprocessStatus::Starting)
397                        || r.error != before_error
398                        || r.summary != before_summary
399                        || r.desired_enabled != before_desired_enabled
400                })
401            };
402            changed |= changed_this_record;
403
404            self.refresh_background_archive_metadata(&record_id).await?;
405        }
406
407        if changed {
408            self.save_background_state().await?;
409        }
410        Ok(self.background_status_entries().await)
411    }
412
413    async fn update_background_record_state(
414        &self,
415        record_id: &str,
416        snapshot: Option<crate::tools::types::VTCodeExecSession>,
417    ) -> Result<Option<(String, String, u8)>> {
418        let termination_requested = if let Some(snapshot) = snapshot.as_ref() {
419            self.config.exec_sessions.termination_requested(snapshot.id.as_str()).await
420        } else {
421            false
422        };
423        let mut state = self.state.write().await;
424        let Some(record) = state.background_children.get_mut(record_id) else {
425            return Ok(None);
426        };
427        let Some(snapshot) = snapshot else {
428            return Self::handle_missing_background_snapshot(record, &self.config);
429        };
430
431        // A previous execution may finish while the stable task is restarting.
432        if record.exec_session_id != snapshot.id.as_str() {
433            return Ok(None);
434        }
435        record.updated_at = Utc::now();
436        record.exit_code = snapshot.exit_code;
437        record.termination_requested |= termination_requested;
438        record.pid = snapshot.child_pid;
439        record.started_at = snapshot.started_at.or(record.started_at);
440
441        match snapshot.lifecycle_state {
442            Some(crate::tools::types::VTCodeSessionLifecycleState::Running) => {
443                // A graceful stop sets `desired_enabled=false` optimistically
444                // while SIGTERM drains. Do not resurrect `Stopped` back to
445                // `Running` during that grace window; the `Exited` arm below
446                // finalizes once the backend confirms exit.
447                if !record.desired_enabled && matches!(record.status, BackgroundSubprocessStatus::Stopped) {
448                    return Ok(None);
449                }
450                record.status = BackgroundSubprocessStatus::Running;
451                record.ended_at = None;
452                record.error = None;
453            }
454            Some(crate::tools::types::VTCodeSessionLifecycleState::Exited) | None => {
455                record.ended_at.get_or_insert(Utc::now());
456                // A clean `exit 0` is successful completion, not a crash.
457                // It must not trigger auto-restore and must surface as
458                // `Stopped` so the runloop/drawer agree with the exec
459                // session's `exited (0)` status.
460                if matches!(snapshot.exit_code, Some(0)) {
461                    record.desired_enabled = false;
462                    record.status = BackgroundSubprocessStatus::Stopped;
463                    record.summary = Some("Background subprocess completed successfully".to_string());
464                    record.error = None;
465                    return Ok(None);
466                }
467                if record.desired_enabled
468                    && self.config.vt_cfg.subagents.background.auto_restore
469                    && record.restart_attempts < 1
470                {
471                    let next_restart_attempt = record.restart_attempts.saturating_add(1);
472                    record.restart_attempts = next_restart_attempt;
473                    record.status = BackgroundSubprocessStatus::Starting;
474                    tracing::warn!(
475                        agent_name = record.agent_name.as_str(),
476                        record_id = record.id.as_str(),
477                        attempt = next_restart_attempt,
478                        "Background subprocess exited unexpectedly; scheduling restart"
479                    );
480                    return Ok(Some((record.agent_name.clone(), record.id.clone(), next_restart_attempt)));
481                }
482                Self::mark_background_record_stopped_or_error(record, &snapshot, &self.config);
483            }
484        }
485
486        Ok(None)
487    }
488
489    fn handle_missing_background_snapshot(
490        record: &mut BackgroundRecord,
491        config: &SubagentControllerConfig,
492    ) -> Result<Option<(String, String, u8)>> {
493        if record.desired_enabled && config.vt_cfg.subagents.background.auto_restore {
494            if record.restart_attempts < 1 {
495                let next_restart_attempt = record.restart_attempts.saturating_add(1);
496                record.restart_attempts = next_restart_attempt;
497                record.status = BackgroundSubprocessStatus::Starting;
498                tracing::warn!(
499                    agent_name = record.agent_name.as_str(),
500                    record_id = record.id.as_str(),
501                    attempt = next_restart_attempt,
502                    "Background subprocess is missing; scheduling restart"
503                );
504                return Ok(Some((record.agent_name.clone(), record.id.clone(), next_restart_attempt)));
505            }
506            record.status = BackgroundSubprocessStatus::Error;
507            record.error = Some("Background subprocess is not running".to_string());
508            record.ended_at.get_or_insert(Utc::now());
509        } else if !record.desired_enabled {
510            record.status = BackgroundSubprocessStatus::Stopped;
511            record.ended_at.get_or_insert(Utc::now());
512        }
513        Ok(None)
514    }
515
516    fn mark_background_record_stopped_or_error(
517        record: &mut BackgroundRecord,
518        snapshot: &crate::tools::types::VTCodeExecSession,
519        _config: &SubagentControllerConfig,
520    ) {
521        // Defensive: a retained `Some(0)` snapshot must never surface as
522        // `Error`. This covers paths that bypass the early clean-exit return
523        // above (e.g. restart budget already exhausted).
524        if matches!(snapshot.exit_code, Some(0)) {
525            record.desired_enabled = false;
526            record.status = BackgroundSubprocessStatus::Stopped;
527            record.summary = Some("Background subprocess completed successfully".to_string());
528            record.error = None;
529            record.ended_at.get_or_insert(Utc::now());
530            return;
531        }
532        if record.desired_enabled {
533            record.status = BackgroundSubprocessStatus::Error;
534            record.summary = None;
535            record.error = Some(match snapshot.exit_code {
536                Some(exit_code) => format!("Background subprocess exited with code {exit_code}"),
537                None => "Background subprocess exited unexpectedly".to_string(),
538            });
539        } else {
540            record.status = BackgroundSubprocessStatus::Stopped;
541            record.summary = Some("Background subprocess stopped".to_string());
542            record.error = None;
543        }
544    }
545
546    /// Blocks until one of the target background subprocesses reaches a
547    /// terminal state (`Stopped`/`Error`) or the timeout expires.
548    ///
549    /// This is the background counterpart to the delegated
550    /// [`SubagentController::wait`]: managed subprocesses previously had no
551    /// model-visible wait path, so the main orchestrator could only observe
552    /// completion via manual `/subprocesses` polling or the Local Agents
553    /// drawer. Unknown ids resolve to `Ok(None)` (fail-closed) rather than
554    /// an error so the unified `agent action=wait` dispatcher can race this
555    /// alongside the delegated wait without hallucinating completion.
556    pub async fn wait_for_background(
557        &self,
558        targets: &[String],
559        timeout_ms: Option<u64>,
560    ) -> Result<Option<BackgroundSubprocessEntry>> {
561        if targets.is_empty() {
562            return Ok(None);
563        }
564        let mut completion_rx = self.subscribe_background_completions();
565        let _ = self.refresh_background_processes().await?;
566        for target in targets {
567            if let Ok(entry) = self.background_status_for(target).await
568                && matches!(entry.status, BackgroundSubprocessStatus::Stopped | BackgroundSubprocessStatus::Error)
569            {
570                return Ok(Some(entry));
571            }
572        }
573        let known = {
574            let state = self.state.read().await;
575            targets.iter().any(|target| state.background_children.contains_key(target))
576        };
577        if !known {
578            return Ok(None);
579        }
580
581        let timeout = std::time::Duration::from_millis(
582            timeout_ms.unwrap_or_else(|| self.config.vt_cfg.subagents.default_timeout_seconds.saturating_mul(1000)),
583        );
584        let deadline = tokio::time::Instant::now() + timeout;
585        loop {
586            let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
587            if remaining.is_zero() {
588                return Ok(None);
589            }
590            tokio::select! {
591                result = completion_rx.recv() => {
592                    match result {
593                        Ok(event) if targets.iter().any(|target| target == &event.task_id || target == &event.exec_session_id) => {
594                            for target in targets {
595                                if let Ok(entry) = self.background_status_for(target).await
596                                    && matches!(entry.status, BackgroundSubprocessStatus::Stopped | BackgroundSubprocessStatus::Error)
597                                {
598                                    return Ok(Some(entry));
599                                }
600                            }
601                        }
602                        Ok(_) => {}
603                        Err(broadcast::error::RecvError::Lagged(_)) => {
604                            let _ = self.refresh_background_processes().await?;
605                            for target in targets {
606                                if let Ok(entry) = self.background_status_for(target).await
607                                    && matches!(entry.status, BackgroundSubprocessStatus::Stopped | BackgroundSubprocessStatus::Error)
608                                {
609                                    return Ok(Some(entry));
610                                }
611                            }
612                        }
613                        Err(broadcast::error::RecvError::Closed) => return Ok(None),
614                    }
615                }
616                _ = tokio::time::sleep(remaining) => {
617                    let _ = self.refresh_background_processes().await?;
618                    for target in targets {
619                        if let Ok(entry) = self.background_status_for(target).await
620                            && matches!(entry.status, BackgroundSubprocessStatus::Stopped | BackgroundSubprocessStatus::Error)
621                        {
622                            return Ok(Some(entry));
623                        }
624                    }
625                    return Ok(None);
626                }
627            }
628        }
629    }
630
631    /// Gracefully stops a background subprocess by setting its desired state to disabled.
632    pub async fn graceful_stop_background(&self, target: &str) -> Result<BackgroundSubprocessEntry> {
633        let (agent_name, exec_session_id) = {
634            let mut state = self.state.write().await;
635            let record = state
636                .background_children
637                .get_mut(target)
638                .ok_or_else(|| anyhow!("Unknown background subprocess {target}"))?;
639            record.termination_requested = true;
640            record.desired_enabled = false;
641            record.status = BackgroundSubprocessStatus::Stopped;
642            record.summary = Some("Background subprocess stopped".to_string());
643            record.error = None;
644            record.updated_at = Utc::now();
645            record.ended_at = Some(Utc::now());
646            (record.agent_name.clone(), record.exec_session_id.clone())
647        };
648
649        tracing::info!(
650            agent_name = agent_name.as_str(),
651            record_id = target,
652            exec_session_id = exec_session_id.as_str(),
653            "Gracefully stopping background subagent subprocess"
654        );
655
656        if !exec_session_id.is_empty() {
657            let _ = self.config.exec_sessions.terminate_session(&exec_session_id).await;
658            let _ = self.config.exec_sessions.prune_exited_session(&exec_session_id).await;
659        }
660
661        self.refresh_background_archive_metadata(target).await?;
662        self.save_background_state().await?;
663        self.background_status_for(target).await
664    }
665
666    /// Force-cancels a background subprocess, closing its exec session immediately.
667    pub async fn force_cancel_background(&self, target: &str) -> Result<BackgroundSubprocessEntry> {
668        let (agent_name, exec_session_id) = {
669            let mut state = self.state.write().await;
670            let record = state
671                .background_children
672                .get_mut(target)
673                .ok_or_else(|| anyhow!("Unknown background subprocess {target}"))?;
674            record.termination_requested = true;
675            record.desired_enabled = false;
676            record.status = BackgroundSubprocessStatus::Stopped;
677            record.summary = Some("Background subprocess stopped".to_string());
678            record.error = None;
679            record.updated_at = Utc::now();
680            record.ended_at = Some(Utc::now());
681            (record.agent_name.clone(), record.exec_session_id.clone())
682        };
683
684        tracing::info!(
685            agent_name = agent_name.as_str(),
686            record_id = target,
687            exec_session_id = exec_session_id.as_str(),
688            "Force cancelling background subagent subprocess"
689        );
690
691        if !exec_session_id.is_empty() {
692            let _ = self.config.exec_sessions.close_session(&exec_session_id).await;
693        }
694
695        self.refresh_background_archive_metadata(target).await?;
696        self.save_background_state().await?;
697        self.publish_terminal_background_completion(target, &exec_session_id, None)
698            .await?;
699        self.background_status_for(target).await
700    }
701
702    /// Returns a thread snapshot for a tracked child subagent by its target id.
703    pub async fn snapshot_for_thread(&self, target: &str) -> Result<SubagentThreadSnapshot> {
704        let (
705            id,
706            session_id,
707            parent_thread_id,
708            agent_name,
709            display_label,
710            status,
711            background,
712            created_at,
713            updated_at,
714            archive_path,
715            transcript_path,
716            effective_config,
717            thread_handle,
718            archive_metadata,
719            stored_messages,
720            recent_events,
721        ) = {
722            let state = self.state.read().await;
723            let record = state
724                .children
725                .get(target)
726                .ok_or_else(|| anyhow!("Unknown subagent id {target}"))?;
727            (
728                record.id.clone(),
729                record.session_id.clone(),
730                record.parent_thread_id.clone(),
731                record.spec.name.clone(),
732                record.display_label.clone(),
733                record.status,
734                record.background,
735                record.created_at,
736                record.updated_at,
737                record.archive_path.clone(),
738                record.transcript_path.clone(),
739                record.effective_config.clone(),
740                record.thread_handle.clone(),
741                record.archive_metadata.clone(),
742                record.stored_messages.clone(),
743                record
744                    .thread_handle
745                    .as_ref()
746                    .map(ThreadRuntimeHandle::recent_events)
747                    .unwrap_or_default(),
748            )
749        };
750
751        let effective_config = effective_config
752            .ok_or_else(|| anyhow!("Subagent {target} does not have a captured runtime configuration yet"))?;
753        let snapshot = match thread_handle {
754            Some(handle) => handle.snapshot(),
755            None => {
756                let archive_listing = match archive_path.as_ref() {
757                    Some(path) if tokio::fs::metadata(path).await.is_ok() => load_session_listing(path).await.ok(),
758                    _ => None,
759                };
760                let metadata = archive_listing
761                    .as_ref()
762                    .map(|listing| listing.snapshot.metadata.clone())
763                    .or(archive_metadata)
764                    .or_else(|| {
765                        Some(crate::core::threads::build_thread_archive_metadata(
766                            &self.config.workspace_root,
767                            effective_config.agent.default_model.as_str(),
768                            effective_config.agent.provider.as_str(),
769                            effective_config.agent.theme.as_str(),
770                            effective_config.agent.reasoning_effort.as_str(),
771                        ))
772                    });
773                ThreadSnapshot {
774                    thread_id: ThreadId::new(session_id.clone()),
775                    metadata,
776                    archive_listing,
777                    messages: stored_messages,
778                    loaded_skills: Vec::new(),
779                    turn_in_flight: false,
780                }
781            }
782        };
783
784        Ok(SubagentThreadSnapshot {
785            id,
786            session_id,
787            parent_thread_id,
788            agent_name,
789            display_label,
790            status,
791            background,
792            created_at,
793            updated_at,
794            archive_path,
795            transcript_path,
796            effective_config,
797            snapshot,
798            recent_events,
799        })
800    }
801}
802
803#[cfg(test)]
804mod tests {
805    use super::*;
806    use crate::subagents::tests::{read_only_test_spec, test_background_record, test_controller_config};
807
808    #[tokio::test]
809    async fn program_status_stale_snapshot_does_not_settle_restarted_background_task() {
810        let temp = tempfile::TempDir::new().unwrap();
811        let controller =
812            SubagentController::new(test_controller_config(temp.path().to_path_buf(), VTCodeConfig::default()))
813                .await
814                .unwrap();
815        let spec = read_only_test_spec("demo");
816        let record =
817            test_background_record(&spec, "stable-task", BackgroundSubprocessStatus::Running, true, "exec-current");
818        let updated_at = record.updated_at;
819        controller
820            .state
821            .write()
822            .await
823            .background_children
824            .insert(record.id.clone(), record);
825        let snapshot = serde_json::from_value(serde_json::json!({
826            "id":"exec-old", "backend":"pipe", "command":"ignored", "args":[], "lifecycle_state":"exited", "exit_code":19
827        })).unwrap();
828        assert!(
829            controller
830                .update_background_record_state("stable-task", Some(snapshot))
831                .await
832                .unwrap()
833                .is_none()
834        );
835        let state = controller.state.read().await;
836        let record = &state.background_children["stable-task"];
837        assert_eq!(record.status, BackgroundSubprocessStatus::Running);
838        assert_eq!(record.exec_session_id, "exec-current");
839        assert_eq!(record.exit_code, None);
840        assert_eq!(record.updated_at, updated_at);
841        assert!(!record.termination_requested);
842    }
843
844    #[test]
845    fn program_status_completion_evidence_survives_persistence() {
846        let spec = read_only_test_spec("demo");
847        let mut record =
848            test_background_record(&spec, "stable-task", BackgroundSubprocessStatus::Stopped, false, "exec-terminated");
849        record.exit_code = Some(137);
850        record.termination_requested = true;
851        let persisted = record.into_persisted();
852        let bytes = serde_json::to_vec(&persisted).unwrap();
853        let decoded: PersistedBackgroundRecord = serde_json::from_slice(&bytes).unwrap();
854        let restored = BackgroundRecord::from_persisted(decoded).build_status_entry();
855        assert_eq!(restored.exit_code, Some(137));
856        assert!(restored.termination_requested);
857    }
858}