1use super::*;
4
5#[derive(Default)]
6pub(crate) struct SessionPreviewState {
7 pub(crate) retained: RetainedSessionPreview,
8 pub(crate) pending: Option<String>,
9}
10
11#[derive(Clone)]
12pub struct ExecSessionManager {
13 pipe_sessions: PipeSessionManager,
14 pty_sessions: PtySessionManager,
15 sessions: Arc<RwLock<HashMap<ExecSessionId, Arc<ExecSessionRecord>>>>,
16 create_lock: Arc<Mutex<()>>,
17 active_background_processes: Arc<AtomicUsize>,
18 pub(crate) foreground_pty_counter: Arc<ParkingMutex<Option<Arc<AtomicUsize>>>>,
19 pub(crate) foreground_session: Arc<ParkingMutex<Option<ExecSessionId>>>,
20 focused_session: Arc<ParkingMutex<Option<ExecSessionId>>>,
21 background_request: Arc<ParkingMutex<Option<ExecSessionId>>>,
22 background_shortcut_result: Arc<ParkingMutex<Option<BackgroundShortcutResult>>>,
23 completion_tx: broadcast::Sender<ExecSessionCompletionEvent>,
24 completion_notify: Arc<Notify>,
25}
26
27impl ExecSessionManager {
28 #[must_use]
29 pub fn new(workspace_root: PathBuf, pty_sessions: PtySessionManager) -> Self {
30 let (completion_tx, _) = broadcast::channel(64);
31 Self {
32 pipe_sessions: PipeSessionManager::new(workspace_root),
33 pty_sessions,
34 sessions: Arc::new(RwLock::new(HashMap::new())),
35 create_lock: Arc::new(Mutex::new(())),
36 active_background_processes: Arc::new(AtomicUsize::new(0)),
37 foreground_pty_counter: Arc::new(ParkingMutex::new(None)),
38 foreground_session: Arc::new(ParkingMutex::new(None)),
39 focused_session: Arc::new(ParkingMutex::new(None)),
40 background_request: Arc::new(ParkingMutex::new(None)),
41 background_shortcut_result: Arc::new(ParkingMutex::new(None)),
42 completion_tx,
43 completion_notify: Arc::new(Notify::new()),
44 }
45 }
46
47 pub fn subscribe_completion(&self) -> broadcast::Receiver<ExecSessionCompletionEvent> {
49 self.completion_tx.subscribe()
50 }
51
52 #[must_use]
54 pub fn completion_notify(&self) -> Arc<Notify> {
55 Arc::clone(&self.completion_notify)
56 }
57
58 pub(crate) fn set_foreground_pty_counter(&self, counter: Arc<AtomicUsize>) {
59 *self.foreground_pty_counter.lock() = Some(counter);
60 }
61
62 #[cfg(test)]
65 pub(crate) async fn create_pipe_session(
66 &self,
67 session_id: ExecSessionId,
68 command: Vec<String>,
69 working_dir: PathBuf,
70 env: HashMap<String, String>,
71 ) -> Result<VTCodeExecSession> {
72 self.create_pipe_session_with_sandbox_and_background(session_id, command, working_dir, env, false, false)
73 .await
74 }
75
76 #[cfg(test)]
77 pub(crate) async fn create_pipe_session_with_sandbox_and_background(
78 &self,
79 session_id: ExecSessionId,
80 command: Vec<String>,
81 working_dir: PathBuf,
82 env: HashMap<String, String>,
83 sandbox_active: bool,
84 background: bool,
85 ) -> Result<VTCodeExecSession> {
86 self.create_pipe_session_with_stdin(
87 session_id,
88 command,
89 working_dir,
90 env,
91 sandbox_active,
92 background,
93 PipeStdinMode::Piped,
94 )
95 .await
96 }
97
98 #[allow(
99 clippy::too_many_arguments,
100 reason = "Launch options preserve the existing internal session API while public pipe runs select stdin explicitly."
101 )]
102 pub(crate) async fn create_pipe_session_with_stdin(
103 &self,
104 session_id: ExecSessionId,
105 command: Vec<String>,
106 working_dir: PathBuf,
107 env: HashMap<String, String>,
108 sandbox_active: bool,
109 background: bool,
110 stdin_mode: PipeStdinMode,
111 ) -> Result<VTCodeExecSession> {
112 let launch_mode = if background {
113 ExecSessionLaunchMode::UserBackground
114 } else {
115 ExecSessionLaunchMode::Foreground
116 };
117 self.create_pipe_session_with_launch_mode(
118 session_id,
119 command,
120 working_dir,
121 env,
122 sandbox_active,
123 launch_mode,
124 stdin_mode,
125 )
126 .await
127 }
128
129 pub(crate) async fn create_pipe_session_for_managed_background(
130 &self,
131 session_id: ExecSessionId,
132 command: Vec<String>,
133 working_dir: PathBuf,
134 env: HashMap<String, String>,
135 ) -> Result<VTCodeExecSession> {
136 self.create_pipe_session_with_launch_mode(
137 session_id,
138 command,
139 working_dir,
140 env,
141 false,
142 ExecSessionLaunchMode::ManagedBackground,
143 PipeStdinMode::Null,
144 )
145 .await
146 }
147
148 #[allow(
149 clippy::too_many_arguments,
150 reason = "Internal session launch carries sandbox, lifecycle, and stdin policy through one spawning boundary."
151 )]
152 async fn create_pipe_session_with_launch_mode(
153 &self,
154 session_id: ExecSessionId,
155 command: Vec<String>,
156 working_dir: PathBuf,
157 env: HashMap<String, String>,
158 sandbox_active: bool,
159 launch_mode: ExecSessionLaunchMode,
160 stdin_mode: PipeStdinMode,
161 ) -> Result<VTCodeExecSession> {
162 let _create_guard = self.create_lock.lock().await;
163 self.ensure_session_absent(&session_id).await?;
164 let slot_reserved = if launch_mode.reserves_background_slot() {
165 self.reserve_background_slot()?
166 } else {
167 false
168 };
169 let env = if sandbox_active {
170 build_sanitized_env(&env, true, false, "exec-session", &[])
171 } else {
172 env
173 };
174 let metadata = match self
175 .pipe_sessions
176 .create_session(session_id.clone(), command, working_dir, env, launch_mode.is_background(), stdin_mode)
177 .await
178 {
179 Ok(metadata) => metadata,
180 Err(error) => {
181 self.release_reserved_background_slot(slot_reserved);
182 return Err(error);
183 }
184 };
185 let record = match self
186 .insert_session(metadata.clone(), ExecSessionBackend::Pipe, None, launch_mode, slot_reserved)
187 .await
188 {
189 Ok(record) => record,
190 Err(error) => {
191 let _ = self.pipe_sessions.close_session(session_id.as_str()).await;
192 self.release_reserved_background_slot(slot_reserved);
193 return Err(error);
194 }
195 };
196 if launch_mode.is_background() {
197 self.start_background_watcher(record, session_id.to_string());
198 } else if launch_mode.sets_foreground_session() {
199 self.set_foreground_session(metadata.id.clone());
200 self.start_foreground_watcher(record, session_id.to_string());
201 }
202 Ok(metadata)
203 }
204
205 #[cfg(test)]
208 pub(crate) async fn create_pty_session(
209 &self,
210 session_id: ExecSessionId,
211 command: Vec<String>,
212 working_dir: PathBuf,
213 size: PtySize,
214 extra_env: HashMap<String, String>,
215 zsh_exec_bridge: Option<ZshExecBridgeSession>,
216 ) -> Result<VTCodeExecSession> {
217 self.create_pty_session_with_sandbox_and_background(
218 session_id,
219 command,
220 working_dir,
221 size,
222 extra_env,
223 zsh_exec_bridge,
224 HashMap::new(),
225 false,
226 false,
227 )
228 .await
229 }
230
231 #[allow(
232 clippy::too_many_arguments,
233 reason = "The constructor keeps sandbox, bridge, and background launch settings explicit at the session boundary."
234 )]
235 pub(crate) async fn create_pty_session_with_sandbox_and_background(
236 &self,
237 session_id: ExecSessionId,
238 command: Vec<String>,
239 working_dir: PathBuf,
240 size: PtySize,
241 extra_env: HashMap<String, String>,
242 zsh_exec_bridge: Option<ZshExecBridgeSession>,
243 trusted_env: HashMap<String, String>,
244 sandbox_active: bool,
245 background: bool,
246 ) -> Result<VTCodeExecSession> {
247 let launch_mode = if background {
248 ExecSessionLaunchMode::UserBackground
249 } else {
250 ExecSessionLaunchMode::Foreground
251 };
252 self.create_pty_session_with_launch_mode(
253 session_id,
254 command,
255 working_dir,
256 size,
257 extra_env,
258 zsh_exec_bridge,
259 trusted_env,
260 sandbox_active,
261 launch_mode,
262 )
263 .await
264 }
265
266 #[allow(
267 clippy::too_many_arguments,
268 reason = "The managed background constructor keeps PTY launch settings explicit at the session boundary."
269 )]
270 pub(crate) async fn create_pty_session_for_managed_background(
271 &self,
272 session_id: ExecSessionId,
273 command: Vec<String>,
274 working_dir: PathBuf,
275 size: PtySize,
276 extra_env: HashMap<String, String>,
277 zsh_exec_bridge: Option<ZshExecBridgeSession>,
278 trusted_env: HashMap<String, String>,
279 sandbox_active: bool,
280 ) -> Result<VTCodeExecSession> {
281 self.create_pty_session_with_launch_mode(
282 session_id,
283 command,
284 working_dir,
285 size,
286 extra_env,
287 zsh_exec_bridge,
288 trusted_env,
289 sandbox_active,
290 ExecSessionLaunchMode::ManagedBackground,
291 )
292 .await
293 }
294
295 #[allow(
296 clippy::too_many_arguments,
297 reason = "The constructor keeps sandbox, bridge, and background launch settings explicit at the session boundary."
298 )]
299 async fn create_pty_session_with_launch_mode(
300 &self,
301 session_id: ExecSessionId,
302 command: Vec<String>,
303 working_dir: PathBuf,
304 size: PtySize,
305 extra_env: HashMap<String, String>,
306 zsh_exec_bridge: Option<ZshExecBridgeSession>,
307 trusted_env: HashMap<String, String>,
308 sandbox_active: bool,
309 launch_mode: ExecSessionLaunchMode,
310 ) -> Result<VTCodeExecSession> {
311 let _create_guard = self.create_lock.lock().await;
312 self.ensure_session_absent(&session_id).await?;
313 let slot_reserved = if launch_mode.reserves_background_slot() {
314 self.reserve_background_slot()?
315 } else {
316 false
317 };
318 let pty_guard = match self.pty_sessions.start_session() {
319 Ok(guard) => guard,
320 Err(error) => {
321 self.release_reserved_background_slot(slot_reserved);
322 return Err(error);
323 }
324 };
325 let metadata = match self.pty_sessions.manager().create_session_with_bridge_sandboxed(
326 session_id.clone().into(),
327 command,
328 working_dir,
329 size,
330 extra_env,
331 zsh_exec_bridge,
332 trusted_env,
333 sandbox_active,
334 ) {
335 Ok(metadata) => metadata,
336 Err(error) => {
337 self.release_reserved_background_slot(slot_reserved);
338 return Err(error);
339 }
340 };
341 let mut exec_metadata = VTCodeExecSession::from(metadata);
342 exec_metadata.background = launch_mode.is_background();
343 let record = match self
344 .insert_session(exec_metadata.clone(), ExecSessionBackend::Pty, Some(pty_guard), launch_mode, slot_reserved)
345 .await
346 {
347 Ok(record) => record,
348 Err(error) => {
349 let _ = self.pty_sessions.manager().close_session(session_id.as_str());
350 self.release_reserved_background_slot(slot_reserved);
351 return Err(error);
352 }
353 };
354 if launch_mode.is_background() {
355 self.start_background_watcher(record, session_id.to_string());
356 } else if launch_mode.sets_foreground_session() {
357 self.set_foreground_session(exec_metadata.id.clone());
358 self.start_foreground_watcher(record, session_id.to_string());
359 }
360 Ok(exec_metadata)
361 }
362
363 pub(crate) async fn snapshot_session(&self, session_id: &str) -> Result<VTCodeExecSession> {
364 let record = self.session_record(session_id).await?;
365 match record.backend {
366 ExecSessionBackend::Pipe => self.pipe_sessions.session_record(session_id).await.map(|r| {
367 let mut metadata = r.metadata.clone();
368 metadata.background = record.background.load(Ordering::Acquire);
369 let exit_code = if r.handle.has_exited() {
370 r.handle.exit_code()
371 } else {
372 None
373 };
374 metadata.exit_code = exit_code;
375 metadata.lifecycle_state = Some(if exit_code.is_some() {
376 crate::tools::types::VTCodeSessionLifecycleState::Exited
377 } else {
378 crate::tools::types::VTCodeSessionLifecycleState::Running
379 });
380 metadata
381 }),
382 ExecSessionBackend::Pty => self.pty_sessions.manager().snapshot_session(session_id).map(|metadata| {
383 let mut metadata = VTCodeExecSession::from(metadata);
384 metadata.background = record.background.load(Ordering::Acquire);
385 metadata
386 }),
387 }
388 }
389
390 pub(crate) async fn termination_requested(&self, session_id: &str) -> bool {
391 self.session_record(session_id)
392 .await
393 .is_ok_and(|record| record.termination_requested.load(Ordering::Acquire))
394 }
395
396 pub async fn background_session_snapshot(&self, session_id: &str) -> Result<ExecSessionUiSnapshot> {
398 let record = self.session_record(session_id).await?;
399 if !record.background.load(Ordering::Acquire) || !record.show_in_background_drawer.load(Ordering::Acquire) {
400 bail!("exec session '{session_id}' is a foreground session and is not visible in the background drawer");
401 }
402
403 let _ = self.read_session_output(session_id, false).await?;
407 let metadata = self.snapshot_session(session_id).await?;
408 let updated_at = {
409 let mut completed_at = record.completed_at.lock();
410 if metadata.exit_code.is_some() {
411 completed_at.get_or_insert_with(Utc::now);
412 }
413 (*completed_at).or(metadata.started_at).unwrap_or_else(Utc::now)
414 };
415 Ok(ExecSessionUiSnapshot {
416 updated_at,
417 metadata,
418 preview: record.preview(),
419 termination_requested: record.termination_requested.load(Ordering::Acquire),
420 })
421 }
422
423 pub async fn background_session_snapshots(&self) -> Vec<ExecSessionUiSnapshot> {
426 let ids = {
427 let sessions = self.sessions.read().await;
428 sessions
429 .values()
430 .filter(|record| {
431 record.background.load(Ordering::Acquire)
432 && record.show_in_background_drawer.load(Ordering::Acquire)
433 })
434 .map(|record| record.metadata.id.clone())
435 .collect::<Vec<_>>()
436 };
437
438 let mut snapshots = Vec::new();
439 for id in ids {
440 if let Ok(snapshot) = self.background_session_snapshot(id.as_str()).await {
441 snapshots.push(snapshot);
442 }
443 }
444 snapshots.sort_by(|left, right| match (left.metadata.started_at, right.metadata.started_at) {
445 (Some(left), Some(right)) => right.cmp(&left),
446 (Some(_), None) => std::cmp::Ordering::Less,
447 (None, Some(_)) => std::cmp::Ordering::Greater,
448 (None, None) => right.metadata.id.cmp(&left.metadata.id),
449 });
450 snapshots
451 }
452
453 pub(crate) async fn list_sessions(&self) -> Vec<VTCodeExecSession> {
454 let sessions = self.sessions.read().await;
455 let mut listed = sessions
456 .values()
457 .map(|record| {
458 let mut metadata = record.metadata.clone();
459 metadata.background = record.background.load(Ordering::Acquire);
460 metadata
461 })
462 .collect::<Vec<_>>();
463 listed.sort_by(|left, right| left.id.cmp(&right.id));
464 listed
465 }
466
467 pub(crate) async fn in_progress_exec_sessions(&self, cap: usize) -> Vec<VTCodeExecSession> {
477 self.collect_in_progress_exec_sessions(cap, true).await
478 }
479
480 pub(crate) async fn in_progress_foreground_exec_sessions(&self, cap: usize) -> Vec<VTCodeExecSession> {
484 self.collect_in_progress_exec_sessions(cap, false).await
485 }
486
487 async fn collect_in_progress_exec_sessions(&self, cap: usize, include_background: bool) -> Vec<VTCodeExecSession> {
488 if cap == 0 {
489 return Vec::new();
490 }
491 let ids = {
492 let sessions = self.sessions.read().await;
493 sessions.keys().cloned().collect::<Vec<_>>()
494 };
495 let mut in_progress = Vec::new();
496 for id in ids {
497 let Ok(session) = self.snapshot_session(id.as_str()).await else {
498 continue;
499 };
500 if session.exit_code.is_none() && (include_background || !session.background) {
501 in_progress.push(session);
502 }
503 }
504 in_progress.sort_by(|left, right| match (left.started_at, right.started_at) {
505 (Some(left), Some(right)) => right.cmp(&left),
506 (Some(_), None) => std::cmp::Ordering::Less,
507 (None, Some(_)) => std::cmp::Ordering::Greater,
508 (None, None) => right.id.cmp(&left.id),
509 });
510 in_progress.truncate(cap);
511 in_progress
512 }
513
514 pub(crate) async fn read_session_output(&self, session_id: &str, drain: bool) -> Result<Option<String>> {
515 let record = self.session_record(session_id).await?;
516 let _output_read_guard = record.output_read_lock.lock().await;
520 let output = match record.backend {
521 ExecSessionBackend::Pipe => self.pipe_sessions.read_session_output(session_id, drain).await,
522 ExecSessionBackend::Pty => self.pty_sessions.manager().read_session_output(session_id, drain),
523 }?;
524 record.remember_output(output.as_deref(), drain);
525 Ok(output)
526 }
527
528 pub(crate) async fn output_stats(&self, session_id: &str) -> Result<Option<PipeOutputStats>> {
529 let record = self.session_record(session_id).await?;
530 match record.backend {
531 ExecSessionBackend::Pipe => self.pipe_sessions.output_stats(session_id).await.map(Some),
532 ExecSessionBackend::Pty => self.pty_sessions.manager().output_stats(session_id).map(|stats| {
533 stats.map(|stats| PipeOutputStats {
534 total_bytes: stats.total_bytes,
535 truncated: stats.truncated,
536 spool_path: stats.spool_path,
537 spool_available: stats.spool_available,
538 spool_complete: stats.spool_complete,
539 spool_integrity: stats.spool_integrity,
540 })
541 }),
542 }
543 }
544
545 pub async fn send_input_to_session(&self, session_id: &str, data: &[u8], append_newline: bool) -> Result<usize> {
546 let record = self.session_record(session_id).await?;
547 match record.backend {
548 ExecSessionBackend::Pipe => {
549 self.pipe_sessions.send_input_to_session(session_id, data, append_newline).await
550 }
551 ExecSessionBackend::Pty => {
552 self.pty_sessions
553 .manager()
554 .send_input_to_session(session_id, data, append_newline)
555 }
556 }
557 }
558
559 pub async fn is_session_completed(&self, session_id: &str) -> Result<Option<i32>> {
560 let record = self.session_record(session_id).await?;
561 let completed = match record.backend {
562 ExecSessionBackend::Pipe => self.pipe_sessions.is_session_completed(session_id).await,
563 ExecSessionBackend::Pty => self.pty_sessions.manager().is_session_completed(session_id),
564 }?;
565 if completed.is_some() {
566 self.release_pending_background_request(session_id);
567 self.clear_focused_session_if_matches(session_id);
568 }
569 Ok(completed)
570 }
571
572 pub(crate) async fn activity_receiver(&self, session_id: &str) -> Result<Option<watch::Receiver<u64>>> {
573 let record = self.session_record(session_id).await?;
574 match record.backend {
575 ExecSessionBackend::Pipe => self.pipe_sessions.activity_receiver(session_id).await.map(Some),
576 ExecSessionBackend::Pty => Ok(None),
577 }
578 }
579
580 pub(crate) async fn is_output_drained(&self, session_id: &str) -> Result<bool> {
581 let record = self.session_record(session_id).await?;
582 let drained = match record.backend {
583 ExecSessionBackend::Pipe => self.pipe_sessions.is_output_drained(session_id).await,
584 ExecSessionBackend::Pty => self.pty_sessions.manager().is_output_drained(session_id),
585 }?;
586 Ok(drained)
587 }
588
589 pub async fn terminate_session(&self, session_id: &str) -> Result<()> {
590 let record = self.session_record(session_id).await?;
591 self.clear_focused_session_if_matches(session_id);
592 record.termination_requested.store(true, Ordering::Release);
593 let result = match record.backend {
594 ExecSessionBackend::Pipe => self.pipe_sessions.terminate_session(session_id).await,
595 ExecSessionBackend::Pty => {
596 let manager = self.pty_sessions.manager().clone();
597 let id = session_id.to_string();
598 tokio::task::spawn_blocking(move || manager.terminate_session(&id))
599 .await
600 .map_err(|join_error| anyhow!("exec session terminate task failed: {join_error}"))?
601 }
602 };
603 if result.is_err() {
604 record.termination_requested.store(false, Ordering::Release);
605 }
606 result
607 }
608
609 pub async fn force_terminate_session(&self, session_id: &str) -> Result<()> {
610 let record = self.session_record(session_id).await?;
611 self.clear_focused_session_if_matches(session_id);
612 record.termination_requested.store(true, Ordering::Release);
613 let result = match record.backend {
614 ExecSessionBackend::Pipe => self.pipe_sessions.force_terminate_session(session_id).await,
615 ExecSessionBackend::Pty => {
616 let manager = self.pty_sessions.manager().clone();
620 let id = session_id.to_string();
621 tokio::task::spawn_blocking(move || manager.force_terminate_session(&id))
622 .await
623 .map_err(|join_error| anyhow!("exec session force-terminate task failed: {join_error}"))?
624 }
625 };
626 if result.is_err() {
627 record.termination_requested.store(false, Ordering::Release);
628 }
629 result
630 }
631
632 pub async fn close_session(&self, session_id: &str) -> Result<VTCodeExecSession> {
633 self.close_session_with_mode(session_id, PtyCloseMode::Graceful).await
634 }
635
636 pub async fn close_session_with_mode(&self, session_id: &str, mode: PtyCloseMode) -> Result<VTCodeExecSession> {
641 let (record, pending_background_request) = {
645 let _lifecycle_guard = self.create_lock.lock().await;
646 let record = {
647 let mut sessions = self.sessions.write().await;
648 sessions
649 .remove(session_id)
650 .ok_or_else(|| missing_exec_session_error(session_id))?
651 };
652
653 let pending_background_request = self.clear_foreground_and_take_pending_request(session_id);
654 self.clear_focused_session_if_matches(session_id);
655 (record, pending_background_request)
656 };
657
658 let close_result = self.close_session_backend_bounded(session_id, &record, mode).await;
662
663 self.release_foreground_pty_count(&record);
664 if pending_background_request {
665 self.release_reserved_background_slot(true);
666 } else {
667 self.release_background_slot(&record);
668 }
669
670 let mut metadata = close_result?;
671 metadata.background = record.background.load(Ordering::Acquire);
672 Ok(metadata)
673 }
674
675 async fn close_session_backend_bounded(
683 &self,
684 session_id: &str,
685 record: &Arc<ExecSessionRecord>,
686 mode: PtyCloseMode,
687 ) -> Result<VTCodeExecSession> {
688 let background_watch = record.background_watch.lock().take();
689 if let Some(watch) = background_watch {
690 watch.abort();
691 if tokio::time::timeout(EXEC_SESSION_WATCH_ABORT_TIMEOUT, watch).await.is_err() {
692 tracing::warn!(%session_id, "background watcher did not stop within abort timeout");
693 }
694 }
695 let foreground_watch = record.foreground_watch.lock().take();
696 if let Some(watch) = foreground_watch {
697 watch.abort();
698 if tokio::time::timeout(EXEC_SESSION_WATCH_ABORT_TIMEOUT, watch).await.is_err() {
699 tracing::warn!(%session_id, "foreground watcher did not stop within abort timeout");
700 }
701 }
702
703 let _output_read_guard =
709 match tokio::time::timeout(EXEC_SESSION_OUTPUT_READ_LOCK_TIMEOUT, record.output_read_lock.lock()).await {
710 Ok(guard) => Some(guard),
711 Err(_elapsed) => {
712 tracing::warn!(
713 %session_id,
714 "output read lock not available within timeout; closing session anyway"
715 );
716 None
717 }
718 };
719 let metadata = match record.backend {
720 ExecSessionBackend::Pipe => {
721 let pipe_sessions = self.pipe_sessions.clone();
722 let session_id_owned = session_id.to_string();
723 tokio::time::timeout(EXEC_SESSION_CLOSE_TIMEOUT, async move {
724 pipe_sessions.close_session(&session_id_owned).await
725 })
726 .await
727 }
728 ExecSessionBackend::Pty => {
729 let pty_manager = self.pty_sessions.manager().clone();
730 let session_id_owned = session_id.to_string();
731 tokio::time::timeout(EXEC_SESSION_CLOSE_TIMEOUT, async move {
732 tokio::task::spawn_blocking(move || {
733 pty_manager
734 .close_session_with_mode(&session_id_owned, mode)
735 .map(VTCodeExecSession::from)
736 })
737 .await
738 .map_err(|join_error| anyhow!("exec session close task failed: {join_error}"))?
739 })
740 .await
741 }
742 };
743
744 match metadata {
745 Ok(result) => result,
746 Err(_elapsed) => Err(anyhow!(
747 "exec session '{session_id}' close timed out after {}s; record detached and counters released",
748 EXEC_SESSION_CLOSE_TIMEOUT.as_secs()
749 )),
750 }
751 }
752
753 pub async fn force_terminate_or_close(&self, session_id: &str) -> Result<bool> {
756 let completed = self.is_session_completed(session_id).await?.is_some();
757 if completed {
758 self.close_session(session_id).await?;
759 } else {
760 self.force_terminate_session(session_id).await?;
761 }
762 Ok(completed)
763 }
764
765 pub async fn focus_background_session(&self, session_id: &str) -> Result<()> {
768 let record = self.session_record(session_id).await?;
769 if !record.background.load(Ordering::Acquire) {
770 bail!("exec session '{session_id}' is a foreground session and cannot be focused from the drawer");
771 }
772 if self.is_session_completed(session_id).await?.is_some() {
773 bail!("exec session '{session_id}' has already exited and cannot receive input");
774 }
775 *self.focused_session.lock() = Some(record.metadata.id.clone());
776 Ok(())
777 }
778
779 #[must_use]
781 pub fn focused_session_id(&self) -> Option<String> {
782 self.focused_session.lock().as_ref().map(|id| id.as_str().to_string())
783 }
784
785 pub fn clear_focused_session(&self) {
786 *self.focused_session.lock() = None;
787 }
788
789 pub(crate) async fn prune_exited_session(&self, session_id: &str) -> Result<Option<VTCodeExecSession>> {
790 let record = self.session_record(session_id).await?;
791 if self.is_session_completed(session_id).await?.is_some() {
792 let completion_pending = record.background.load(Ordering::Acquire)
798 && record
799 .background_watch
800 .lock()
801 .as_ref()
802 .is_some_and(|watch| !watch.is_finished());
803 if completion_pending {
804 return Ok(None);
805 }
806 return self.close_session(session_id).await.map(Some);
807 }
808 Ok(None)
809 }
810
811 pub(crate) async fn terminate_all_sessions_async(&self) -> Result<()> {
812 self.terminate_all_sessions_with_mode_async(PtyCloseMode::Graceful).await
813 }
814
815 pub(crate) async fn terminate_all_sessions_for_exit_async(&self) -> Result<()> {
820 self.terminate_all_sessions_with_mode_async(PtyCloseMode::Immediate).await
821 }
822
823 async fn terminate_all_sessions_with_mode_async(&self, mode: PtyCloseMode) -> Result<()> {
824 let ids = {
825 let sessions = self.sessions.read().await;
826 sessions.keys().cloned().collect::<Vec<_>>()
827 };
828
829 if ids.is_empty() {
835 return self.pipe_sessions.terminate_all_sessions().await;
836 }
837
838 let results = futures::future::join_all(ids.into_iter().map(|session_id| async move {
839 self.close_session_with_mode(&session_id, mode)
840 .await
841 .map_err(|err| format!("{session_id}: {err}"))
842 }))
843 .await;
844
845 let mut failures: Vec<String> = results.into_iter().filter_map(|r| r.err()).collect();
846
847 if let Err(err) = self.pipe_sessions.terminate_all_sessions().await {
848 failures.push(err.to_string());
849 }
850
851 if failures.is_empty() {
852 Ok(())
853 } else {
854 Err(anyhow!("failed to terminate all exec sessions: {}", failures.join("; ")))
855 }
856 }
857
858 pub(crate) async fn terminate_active_sessions_async(&self) -> Result<()> {
859 self.terminate_active_sessions_with_mode_async(PtyCloseMode::Graceful).await
860 }
861
862 pub(crate) async fn terminate_active_sessions_for_exit_async(&self) -> Result<()> {
865 self.terminate_active_sessions_with_mode_async(PtyCloseMode::Immediate).await
866 }
867
868 async fn terminate_active_sessions_with_mode_async(&self, mode: PtyCloseMode) -> Result<()> {
869 let ids = {
870 let sessions = self.sessions.read().await;
871 sessions
872 .values()
873 .filter(|record| !record.background.load(Ordering::Acquire))
874 .map(|record| record.metadata.id.clone())
875 .collect::<Vec<_>>()
876 };
877
878 if ids.is_empty() {
879 return Ok(());
880 }
881
882 let results: Vec<Result<(), String>> =
885 futures::future::join_all(ids.into_iter().map(|session_id| async move {
886 let should_close = self
887 .session_record(session_id.as_str())
888 .await
889 .map(|record| !record.background.load(Ordering::Acquire))
890 .unwrap_or(false);
891 if !should_close {
892 return Ok(());
893 }
894 self.close_session_with_mode(&session_id, mode)
895 .await
896 .map(|_| ())
897 .map_err(|err| format!("{session_id}: {err}"))
898 }))
899 .await;
900
901 let failures: Vec<String> = results.into_iter().filter_map(|r| r.err()).collect();
902
903 if failures.is_empty() {
904 Ok(())
905 } else {
906 Err(anyhow!("failed to terminate active exec sessions: {}", failures.join("; ")))
907 }
908 }
909
910 pub async fn force_cancel_foreground_sessions(&self) -> (usize, usize, usize) {
919 let ids = {
920 let sessions = self.sessions.read().await;
921 sessions
922 .values()
923 .filter(|record| !record.background.load(Ordering::Acquire))
924 .map(|record| record.metadata.id.clone())
925 .collect::<Vec<_>>()
926 };
927
928 let mut stopped = 0usize;
929 let mut closed = 0usize;
930 let mut failed = 0usize;
931 for session_id in ids {
932 let was_running = self.is_session_completed(session_id.as_str()).await.ok().flatten().is_none();
933 if was_running {
934 if let Err(error) = self.force_terminate_session(session_id.as_str()).await {
935 tracing::warn!(%session_id, %error, "force-cancel terminate failed");
936 }
937 }
938 match self.close_session(session_id.as_str()).await {
939 Ok(_) => {
940 if was_running {
941 stopped += 1;
942 } else {
943 closed += 1;
944 }
945 }
946 Err(error) => {
947 tracing::warn!(%session_id, %error, "force-cancel close failed");
948 failed += 1;
949 }
950 }
951 }
952 (stopped, closed, failed)
953 }
954
955 #[must_use]
957 pub fn active_background_processes(&self) -> usize {
958 self.active_background_processes.load(Ordering::Acquire)
959 }
960
961 pub fn request_foreground_background(&self) -> Option<BackgroundShortcutResult> {
967 let mut request = self.background_request.lock();
968 if request.is_some() {
969 *self.background_shortcut_result.lock() = Some(BackgroundShortcutResult::Requested);
970 return Some(BackgroundShortcutResult::Requested);
971 }
972
973 let Some(session_id) = self.foreground_session.lock().clone() else {
978 *self.background_shortcut_result.lock() = None;
979 return None;
980 };
981
982 let result = match self.reserve_background_slot() {
983 Ok(_) => {
984 *request = Some(session_id);
985 BackgroundShortcutResult::Requested
986 }
987 Err(_) => BackgroundShortcutResult::AtCapacity,
988 };
989 *self.background_shortcut_result.lock() = Some(result);
990 Some(result)
991 }
992
993 pub fn take_background_shortcut_result(&self) -> Option<BackgroundShortcutResult> {
995 self.background_shortcut_result.lock().take()
996 }
997
998 pub(crate) async fn promote_requested_session(&self, session_id: &str) -> Result<bool> {
1000 let _lifecycle_guard = self.create_lock.lock().await;
1001 if !self.take_foreground_promotion_request(session_id) {
1002 return Ok(false);
1003 }
1004
1005 let record = match self.session_record(session_id).await {
1006 Ok(record) => record,
1007 Err(error) => {
1008 self.release_reserved_background_slot(true);
1009 return Err(error);
1010 }
1011 };
1012 if record.background.load(Ordering::Acquire) {
1013 self.release_reserved_background_slot(true);
1014 self.clear_foreground_session(session_id);
1015 return Ok(false);
1016 }
1017 match self.is_session_completed(session_id).await {
1018 Ok(Some(_)) => {
1019 self.release_reserved_background_slot(true);
1020 self.clear_foreground_session(session_id);
1021 return Ok(false);
1022 }
1023 Ok(None) => {}
1024 Err(error) => {
1025 self.release_reserved_background_slot(true);
1026 self.set_foreground_session(ExecSessionId::new(session_id));
1027 return Err(error);
1028 }
1029 }
1030
1031 record.background_slot_reserved.store(true, Ordering::Release);
1032 self.release_foreground_pty_count(&record);
1033 record.background.store(true, Ordering::Release);
1034 record.show_in_background_drawer.store(true, Ordering::Release);
1035 self.clear_foreground_session(session_id);
1036 let promotion_marker = Arc::clone(&record);
1037 self.start_background_watcher(record, session_id.to_string());
1038 promotion_marker
1039 .background_promoted_from_foreground
1040 .store(true, Ordering::Release);
1041 Ok(true)
1042 }
1043
1044 pub(crate) async fn take_foreground_promotion(&self, session_id: &str) -> Result<bool> {
1045 let record = self.session_record(session_id).await?;
1046 Ok(record.background_promoted_from_foreground.swap(false, Ordering::AcqRel))
1047 }
1048
1049 fn reserve_background_slot(&self) -> Result<bool> {
1050 let mut current = self.active_background_processes.load(Ordering::Acquire);
1051 loop {
1052 if current >= MAX_BACKGROUND_PROCESSES {
1053 bail!(
1054 "maximum background process limit reached ({MAX_BACKGROUND_PROCESSES}); wait for or close an existing background session before starting another"
1055 );
1056 }
1057 match self.active_background_processes.compare_exchange(
1058 current,
1059 current + 1,
1060 Ordering::AcqRel,
1061 Ordering::Acquire,
1062 ) {
1063 Ok(_) => return Ok(true),
1064 Err(observed) => current = observed,
1065 }
1066 }
1067 }
1068
1069 fn release_reserved_background_slot(&self, reserved: bool) {
1070 if reserved {
1071 self.active_background_processes.fetch_sub(1, Ordering::AcqRel);
1072 }
1073 }
1074
1075 fn release_background_slot(&self, record: &ExecSessionRecord) {
1076 if record.background_slot_reserved.swap(false, Ordering::AcqRel) {
1077 self.active_background_processes.fetch_sub(1, Ordering::AcqRel);
1078 }
1079 }
1080
1081 fn count_foreground_pty_session(&self, record: &ExecSessionRecord) {
1082 let Some(counter) = self.foreground_pty_counter.lock().as_ref().map(Arc::clone) else {
1083 return;
1084 };
1085 counter.fetch_add(1, Ordering::Relaxed);
1086 *record.foreground_pty_counter.lock() = Some(counter);
1087 }
1088
1089 fn release_foreground_pty_count(&self, record: &ExecSessionRecord) {
1090 let Some(counter) = record.foreground_pty_counter.lock().take() else {
1091 return;
1092 };
1093 let _ = counter.fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| current.checked_sub(1));
1094 }
1095
1096 fn clear_foreground_and_take_pending_request(&self, session_id: &str) -> bool {
1097 let mut request = self.background_request.lock();
1098 let taken = request.as_deref() == Some(session_id) && request.take().is_some();
1099 let mut foreground = self.foreground_session.lock();
1100 if foreground.as_deref() == Some(session_id) {
1101 *foreground = None;
1102 }
1103 drop(foreground);
1104 if taken {
1105 *self.background_shortcut_result.lock() = None;
1106 }
1107 taken
1108 }
1109
1110 fn take_foreground_promotion_request(&self, session_id: &str) -> bool {
1111 let mut request = self.background_request.lock();
1112 if request.as_deref() != Some(session_id) || request.take().is_none() {
1113 return false;
1114 }
1115 let mut foreground = self.foreground_session.lock();
1116 if foreground.as_deref() == Some(session_id) {
1117 *foreground = None;
1118 }
1119 drop(foreground);
1120 *self.background_shortcut_result.lock() = None;
1121 true
1122 }
1123
1124 fn release_pending_background_request(&self, session_id: &str) {
1125 if self.clear_foreground_and_take_pending_request(session_id) {
1126 self.release_reserved_background_slot(true);
1127 }
1128 }
1129
1130 fn start_background_watcher(&self, record: Arc<ExecSessionRecord>, session_id: String) {
1131 let manager = self.clone();
1132 let record_for_task = Arc::clone(&record);
1133 let task = tokio::spawn(async move {
1134 loop {
1135 match manager.is_session_completed(session_id.as_str()).await {
1136 Ok(Some(exit_code)) => {
1137 record_for_task.completed_at.lock().get_or_insert_with(Utc::now);
1138 manager.capture_background_completion_output(session_id.as_str()).await;
1139 if record_for_task
1140 .background_completion_published
1141 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1142 .is_ok()
1143 {
1144 let command = bounded_completion_command(&record_for_task.metadata.command_label());
1145 let _ = manager.completion_tx.send(ExecSessionCompletionEvent {
1146 session_id: ExecSessionId::new(session_id.clone()),
1147 command,
1148 managed_background: !record_for_task.show_in_background_drawer.load(Ordering::Acquire),
1149 termination_requested: record_for_task.termination_requested.load(Ordering::Acquire),
1150 exit_code,
1151 });
1152 manager.completion_notify.notify_one();
1157 }
1158 manager.release_background_slot(&record_for_task);
1159 break;
1160 }
1161 Ok(None) => tokio::time::sleep(tokio::time::Duration::from_millis(50)).await,
1162 Err(_) => {
1163 if manager.session_record(session_id.as_str()).await.is_err() {
1168 break;
1169 }
1170 tokio::time::sleep(tokio::time::Duration::from_millis(50)).await;
1171 }
1172 }
1173 }
1174 });
1175 if let Some(previous) = record.background_watch.lock().replace(task) {
1176 previous.abort();
1177 }
1178 }
1179
1180 async fn capture_background_completion_output(&self, session_id: &str) {
1181 let deadline = tokio::time::Instant::now() + EXEC_SESSION_COMPLETION_DRAIN_TIMEOUT;
1182 loop {
1183 let _ = self.read_session_output(session_id, false).await;
1184 if self.is_output_drained(session_id).await.unwrap_or(false) || tokio::time::Instant::now() >= deadline {
1185 break;
1186 }
1187 tokio::time::sleep(EXEC_SESSION_COMPLETION_DRAIN_POLL).await;
1188 }
1189 let _ = self.read_session_output(session_id, false).await;
1191 }
1192
1193 fn start_foreground_watcher(&self, record: Arc<ExecSessionRecord>, session_id: String) {
1194 let manager = self.clone();
1195 let record_for_task = Arc::clone(&record);
1196 let task = tokio::spawn(async move {
1197 loop {
1198 if manager.promote_requested_session(session_id.as_str()).await.unwrap_or(false) {
1199 break;
1200 }
1201 if record_for_task.background.load(Ordering::Acquire) {
1202 break;
1203 }
1204 match manager.is_session_completed(session_id.as_str()).await {
1205 Ok(Some(_)) => {
1206 manager.release_foreground_pty_count(&record_for_task);
1207 manager.release_pending_background_request(session_id.as_str());
1208 break;
1209 }
1210 Ok(None) => tokio::time::sleep(tokio::time::Duration::from_millis(50)).await,
1211 Err(_) => {
1212 if manager.session_record(session_id.as_str()).await.is_err() {
1213 break;
1214 }
1215 tokio::time::sleep(tokio::time::Duration::from_millis(50)).await;
1216 }
1217 }
1218 }
1219 });
1220 if let Some(previous) = record.foreground_watch.lock().replace(task) {
1221 previous.abort();
1222 }
1223 }
1224
1225 fn set_foreground_session(&self, session_id: ExecSessionId) {
1226 *self.foreground_session.lock() = Some(session_id);
1227 }
1228
1229 fn clear_foreground_session(&self, session_id: &str) {
1230 let mut foreground = self.foreground_session.lock();
1231 if foreground.as_deref() == Some(session_id) {
1232 *foreground = None;
1233 }
1234 }
1235
1236 fn clear_focused_session_if_matches(&self, session_id: &str) {
1237 let mut focused = self.focused_session.lock();
1238 if focused.as_deref().is_some_and(|id| id == session_id) {
1239 *focused = None;
1240 }
1241 }
1242
1243 async fn insert_session(
1244 &self,
1245 metadata: VTCodeExecSession,
1246 backend: ExecSessionBackend,
1247 pty_guard: Option<PtySessionGuard>,
1248 launch_mode: ExecSessionLaunchMode,
1249 background_slot_reserved: bool,
1250 ) -> Result<Arc<ExecSessionRecord>> {
1251 let mut sessions = self.sessions.write().await;
1252 use hashbrown::hash_map::Entry;
1253 match sessions.entry(metadata.id.clone()) {
1254 Entry::Occupied(_) => Err(anyhow!("exec session '{}' already exists", metadata.id.as_str())),
1255 Entry::Vacant(entry) => {
1256 let record = Arc::new(ExecSessionRecord::new(
1257 metadata,
1258 backend,
1259 pty_guard,
1260 launch_mode,
1261 background_slot_reserved,
1262 ));
1263 if launch_mode.sets_foreground_session() {
1266 self.count_foreground_pty_session(&record);
1267 }
1268 entry.insert(Arc::clone(&record));
1269 Ok(record)
1270 }
1271 }
1272 }
1273
1274 async fn ensure_session_absent(&self, session_id: &str) -> Result<()> {
1275 let sessions = self.sessions.read().await;
1276 if sessions.contains_key(session_id) {
1277 return Err(anyhow!("exec session '{session_id}' already exists"));
1278 }
1279 Ok(())
1280 }
1281
1282 pub(crate) async fn session_record(&self, session_id: &str) -> Result<Arc<ExecSessionRecord>> {
1283 let sessions = self.sessions.read().await;
1284 sessions
1285 .get(session_id)
1286 .cloned()
1287 .ok_or_else(|| missing_exec_session_error(session_id))
1288 }
1289}