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 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 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 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 #[must_use]
259 pub fn background_subagents_enabled(&self) -> bool {
260 self.config.vt_cfg.subagents.background.enabled
261 }
262
263 #[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 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 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 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 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 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 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 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 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 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 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 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}