Skip to main content

vtcode_core/tools/exec_session/
manager.rs

1//! `ExecSessionManager`: unified pipe/PTY session management.
2
3use 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    /// Subscribe to confirmed terminal transitions for background exec sessions.
48    pub fn subscribe_completion(&self) -> broadcast::Receiver<ExecSessionCompletionEvent> {
49        self.completion_tx.subscribe()
50    }
51
52    /// Wake an idle interaction loop when a background exec session completes.
53    #[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    /// Test-only convenience: production paths go through
63    /// [`Self::create_pipe_session_with_stdin`].
64    #[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    /// Test-only convenience: production paths go through
206    /// [`Self::create_pty_session_with_sandbox_and_background`].
207    #[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    /// Return one retained background session snapshot for the Local Agents drawer.
397    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        // Capture output that has not yet been consumed by the tool wait loop
404        // before reading the retained preview. Drained output is remembered by
405        // `read_session_output`, so this remains inspectable after a tool turn.
406        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    /// Return all background raw command sessions, including exited sessions
424    /// that have not been explicitly closed.
425    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    /// Bounded snapshot of exec sessions that are still running (not exited).
468    ///
469    /// Used for turn-end diagnostics and telemetry. Ordered newest-first by
470    /// `started_at` (sessions without a timestamp last), capped so a
471    /// pathological session count cannot inflate the recorded state.
472    /// Completion is checked against the backend (not the cached metadata) so
473    /// a session that exited after its last metadata refresh is correctly
474    /// excluded. Cross-turn resume hints use the foreground-only variant
475    /// below so retained background work does not become mandatory follow-up.
476    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    /// Bounded snapshot of running foreground sessions for cross-turn resume
481    /// hints. Retained background sessions are deliberately excluded because
482    /// they are not work the next turn must settle before proceeding.
483    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        // Serialize the backend read with the preview update. A peek followed
517        // by a concurrent drain must not leave the same chunk pending after
518        // it has already been retained.
519        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                // PTY terminate performs a bounded child reap that sleeps;
617                // keep it off the async worker so ForceCancel over N sessions
618                // cannot stall the runloop.
619                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    /// Close with an explicit termination mode. `Immediate` is reserved for
637    /// exit-driven teardown: it skips the `exit\n` courtesy and the SIGTERM
638    /// grace window so a live PTY child cannot cost ~600 ms of shell-return
639    /// latency per session. The OS reaps survivors at process exit.
640    pub async fn close_session_with_mode(&self, session_id: &str, mode: PtyCloseMode) -> Result<VTCodeExecSession> {
641        // Serialize removal with Ctrl+B promotion and new session creation so
642        // a promotion cannot reserve a slot on a record that close has already
643        // detached from the unified session map.
644        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        // Capacity release must happen even if backend close times out: the
659        // record is already detached, so leaving counters elevated pins
660        // `Running PTY command...` and the composer lock forever.
661        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    /// Abort lifecycle watchers and close the backend under a hard timeout.
676    ///
677    /// Watcher abort and the blocking PTY close are themselves the hang sites
678    /// (unbounded `join`/`wait`, or a watcher stuck in a sync section), so
679    /// this entire unit is time-bounded and the PTY close runs on
680    /// `spawn_blocking` to keep the async worker responsive. Pipe close stays
681    /// on the runtime (already async and internally timed).
682    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        // Do not close the backend while an output peek/drain is still using
704        // it. The unified record has already been removed, so this lock only
705        // waits for in-flight readers acquired before close. The acquire is
706        // itself time-bounded: an abandoned peek holding the lock must not
707        // make close (and the runloop) wait forever.
708        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    /// Force-stop an active session, or close it when it has already exited.
754    /// Returns `true` when the session was already complete and was closed.
755    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    /// Focus a running background session so subsequent submitted lines are
766    /// written to its stdin instead of becoming a model prompt.
767    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    /// Return the focused background session, if any.
780    #[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            // Background completion owns the final output capture and parent
793            // notification. A synchronous wait/stop can observe process exit
794            // first; pruning here would abort that watcher and lose the
795            // terminal event. Retain the session until the watcher has
796            // finished, after which a later prune may close it safely.
797            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    /// Exit-path terminator: closes every session in `Immediate` mode so a
816    /// live PTY child is group-SIGKILLed instead of waiting out the SIGTERM
817    /// grace window. Callers still bound this with an outer timeout; the OS
818    /// reaps any remainder at process exit.
819    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        // Parallel close: sequential closes sum per-session timeouts (12s
830        // inner cap each) into a multi-second exit tail. Concurrent closes
831        // overlap as max instead of sum; the caller's outer timeout
832        // (EXIT_BACKGROUND_SHUTDOWN_TIMEOUT) still bounds the total.
833        // Early return when empty avoids even the pipe-manager lock.
834        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    /// Exit-path variant of [`Self::terminate_active_sessions_async`] using
863    /// `Immediate` PTY termination (see [`Self::terminate_all_sessions_for_exit_async`]).
864    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        // Parallel close for the same reason as terminate_all: avoid summing
883        // per-session timeouts on the exit path.
884        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    /// Force-stop every foreground exec session and close it. Used by the TUI
911    /// ForceCancel escape hatch so a stuck PTY cannot keep the composer locked.
912    /// Closing after the kill is intentional: `force_terminate_or_close` alone
913    /// leaves live sessions attached ("remains visible until closed"), which
914    /// would keep the foreground counter elevated. Background sessions are
915    /// user-owned and are left alone. Returns `(stopped, closed, failed)`
916    /// counts where `stopped` are live kills and `closed` are already-exited
917    /// sessions removed from the map.
918    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    /// Return the number of currently live background process reservations.
956    #[must_use]
957    pub fn active_background_processes(&self) -> usize {
958        self.active_background_processes.load(Ordering::Acquire)
959    }
960
961    /// Request that the current foreground session be promoted to background.
962    ///
963    /// The request is synchronous because it is called by the TUI key-event
964    /// callback. The execution wait loop consumes the request asynchronously
965    /// and performs the metadata transition without killing the process.
966    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        // Take the foreground lock after `background_request`. Completion and
974        // close paths use the same order when clearing a pending promotion;
975        // holding both here prevents a completed session from being observed
976        // between the foreground lookup and request reservation.
977        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    /// Consume the result associated with the queued Ctrl+B event.
994    pub fn take_background_shortcut_result(&self) -> Option<BackgroundShortcutResult> {
995        self.background_shortcut_result.lock().take()
996    }
997
998    /// Apply a pending Ctrl+B promotion to the specified session.
999    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                            // Keep one wake permit when the idle loop has not
1153                            // entered its wait yet; the broadcast receiver
1154                            // drains every completion that accumulated behind
1155                            // the single permit.
1156                            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                        // A transient status error must not release the live
1164                        // background reservation. Only a removed session ends
1165                        // this watcher without a confirmed process exit; its
1166                        // close path owns the reservation release.
1167                        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        // Capture bytes that raced with the final drain-state observation.
1190        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                // Foreground PTY and pipe sessions both support Ctrl+B backgrounding,
1264                // so both increment the shared foreground counter driving the TUI hint.
1265                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}