Skip to main content

vtcode_ui/tui/core_tui/runner/
mod.rs

1use std::io;
2use std::time::Duration;
3
4use anyhow::{Context, Result};
5use ratatui::{Terminal, backend::CrosstermBackend};
6use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender};
7use tokio_util::sync::CancellationToken;
8
9use crate::tui::config::types::UiSurfacePreference;
10use crate::tui::options::FullscreenInteractionSettings;
11use crate::tui::ui::tui::log::{clear_tui_log_sender, register_tui_log_sender, set_log_theme_name};
12
13type EventCallback<E> = std::sync::Arc<dyn Fn(&E) + Send + Sync + 'static>;
14
15pub trait TuiCommand {
16    fn is_suspend_event_loop(&self) -> bool;
17    fn is_resume_event_loop(&self) -> bool;
18    fn is_clear_input_queue(&self) -> bool;
19    fn is_force_redraw(&self) -> bool;
20    fn is_stop_event_stream(&self) -> bool;
21    fn is_start_event_stream(&self) -> bool;
22}
23
24pub trait TuiSessionDriver {
25    type Command: TuiCommand;
26    type Event;
27
28    fn handle_command(&mut self, command: Self::Command);
29    #[expect(
30        clippy::type_complexity,
31        reason = "Intentional compatibility, platform, test, or API-shape suppression."
32    )]
33    fn handle_event(
34        &mut self,
35        event: crossterm::event::Event,
36        events: &UnboundedSender<Self::Event>,
37        callback: Option<&(dyn Fn(&Self::Event) + Send + Sync + 'static)>,
38    );
39    fn handle_tick(&mut self);
40    /// Whether the session needs active-rate ticks (animations or transient
41    /// expiries pending). The drive loop uses this to extend active-mode
42    /// ticking independently of recent user input, so truly idle sessions rest
43    /// instead of polling. Input, commands, and crossterm events always render
44    /// immediately regardless of this flag.
45    fn needs_animation_tick(&self) -> bool;
46    fn render(&mut self, frame: &mut ratatui::Frame<'_>);
47    fn take_redraw(&mut self) -> bool;
48    fn use_steady_cursor(&self) -> bool;
49    fn is_hovering_link(&self) -> bool;
50    fn is_selecting_text(&self) -> bool;
51    fn should_exit(&self) -> bool;
52    fn request_exit(&mut self);
53    fn mark_dirty(&mut self);
54    fn update_terminal_title(&mut self);
55    fn attach_program_status_terminal(&mut self);
56    fn flush_program_status(&mut self);
57    fn shutdown_program_status(&mut self);
58    fn clear_terminal_title(&mut self);
59    fn is_running_activity(&self) -> bool;
60    fn has_status_spinner(&self) -> bool;
61    fn thinking_spinner_active(&self) -> bool;
62    fn has_active_navigation_ui(&self) -> bool;
63    fn apply_coalesced_scroll(&mut self, line_delta: i32, page_delta: i32);
64    fn set_show_logs(&mut self, show: bool);
65    fn set_active_pty_sessions(&mut self, sessions: Option<std::sync::Arc<std::sync::atomic::AtomicUsize>>);
66    fn set_workspace_root(&mut self, root: Option<std::path::PathBuf>);
67    fn set_log_receiver(&mut self, receiver: UnboundedReceiver<crate::tui::core_tui::log::LogEntry>);
68    fn set_fullscreen_active(&mut self, active: bool);
69    fn set_fullscreen_interaction(&mut self, config: FullscreenInteractionSettings);
70    fn set_preview_callback(&mut self, callback: Option<PreviewCallback>);
71}
72
73impl TuiCommand for crate::tui::core_tui::types::InlineCommand {
74    fn is_suspend_event_loop(&self) -> bool {
75        matches!(self, crate::tui::core_tui::types::InlineCommand::SuspendEventLoop)
76    }
77
78    fn is_resume_event_loop(&self) -> bool {
79        matches!(self, crate::tui::core_tui::types::InlineCommand::ResumeEventLoop)
80    }
81
82    fn is_clear_input_queue(&self) -> bool {
83        matches!(self, crate::tui::core_tui::types::InlineCommand::ClearInputQueue)
84    }
85
86    fn is_force_redraw(&self) -> bool {
87        matches!(self, crate::tui::core_tui::types::InlineCommand::ForceRedraw)
88    }
89
90    fn is_stop_event_stream(&self) -> bool {
91        matches!(self, crate::tui::core_tui::types::InlineCommand::StopEventStream)
92    }
93
94    fn is_start_event_stream(&self) -> bool {
95        matches!(self, crate::tui::core_tui::types::InlineCommand::StartEventStream)
96    }
97}
98
99use super::types::FocusChangeCallback;
100pub(crate) use super::types::PreviewCallback;
101
102mod drive;
103mod events;
104mod signal;
105mod surface;
106pub(crate) mod terminal_io;
107mod terminal_modes;
108
109use drive::{DriveRuntimeOptions, drive_terminal};
110use events::{EventListener, EventSender, TerminalEvent, spawn_event_loop};
111use signal::SignalCleanupGuard;
112use surface::TerminalSurface;
113use terminal_io::{drain_terminal_events, finalize_terminal, prepare_terminal};
114use terminal_modes::{TerminalModeState, enable_terminal_modes, restore_terminal_modes};
115
116/// Controls the lifecycle of the async crossterm event stream.
117///
118/// The event loop must be fully stopped before launching an external editor
119/// that needs stdin (e.g., nvim), otherwise the background EventStream task
120/// competes with the editor for terminal input, causing freezes.
121pub(super) struct EventStreamController {
122    cancellation_token: CancellationToken,
123    join_handle: Option<tokio::task::JoinHandle<()>>,
124    event_tx: EventSender,
125    rx_paused: std::sync::Arc<std::sync::atomic::AtomicBool>,
126    last_input_elapsed_ms: std::sync::Arc<std::sync::atomic::AtomicU64>,
127    session_start: std::time::Instant,
128}
129
130impl EventStreamController {
131    fn new(
132        cancellation_token: CancellationToken,
133        join_handle: tokio::task::JoinHandle<()>,
134        event_tx: EventSender,
135        rx_paused: std::sync::Arc<std::sync::atomic::AtomicBool>,
136        last_input_elapsed_ms: std::sync::Arc<std::sync::atomic::AtomicU64>,
137        session_start: std::time::Instant,
138    ) -> Self {
139        Self {
140            cancellation_token,
141            join_handle: Some(join_handle),
142            event_tx,
143            rx_paused,
144            last_input_elapsed_ms,
145            session_start,
146        }
147    }
148
149    /// Cancel the current event loop task and await its termination.
150    /// Creates a fresh CancellationToken for the next `start()` call.
151    ///
152    /// Bounded: a stuck crossterm reader must not park TUI teardown, and a
153    /// leaked reader would keep stdin claimed after exit. The handle is
154    /// aborted on timeout so shutdown always makes progress.
155    async fn stop(&mut self) {
156        self.cancellation_token.cancel();
157        if let Some(mut handle) = self.join_handle.take()
158            && tokio::time::timeout(Duration::from_millis(100), &mut handle).await.is_err()
159        {
160            tracing::debug!("event loop did not stop within 100ms; aborting");
161            handle.abort();
162        }
163        self.cancellation_token = CancellationToken::new();
164    }
165
166    /// Spawn a new event loop task with a fresh EventStream.
167    /// Safe to call multiple times after `stop()`.
168    fn start(&mut self) {
169        let token = self.cancellation_token.clone();
170        let event_tx = self.event_tx.clone();
171        let rx_paused = self.rx_paused.clone();
172        let last_input = self.last_input_elapsed_ms.clone();
173        let session_start = self.session_start;
174        self.join_handle = Some(tokio::spawn(async move {
175            spawn_event_loop(event_tx, token, rx_paused, last_input, session_start).await;
176        }));
177    }
178
179    /// Ensure the event loop is stopped for final cleanup on TUI exit.
180    ///
181    /// Bounded like [`Self::stop`]: abort on timeout so a wedged reader can
182    /// never delay the alternate-screen teardown and shell return.
183    async fn shutdown(&mut self) {
184        if let Some(mut handle) = self.join_handle.take() {
185            self.cancellation_token.cancel();
186            if tokio::time::timeout(Duration::from_millis(100), &mut handle).await.is_err() {
187                tracing::debug!("event loop did not shut down within 100ms; aborting");
188                handle.abort();
189            }
190        }
191    }
192}
193
194struct TerminalModeRestoreGuard {
195    state: Option<TerminalModeState>,
196}
197
198impl TerminalModeRestoreGuard {
199    fn new(state: TerminalModeState) -> Self {
200        Self { state: Some(state) }
201    }
202
203    fn state_mut(&mut self) -> &mut TerminalModeState {
204        self.state
205            .as_mut()
206            .expect("terminal mode restore guard must stay armed until shutdown")
207    }
208
209    fn restore(&mut self) -> Result<()> {
210        if let Some(state) = self.state.take() {
211            restore_terminal_modes(&state)?;
212        }
213        Ok(())
214    }
215
216    fn restore_silently(&mut self) {
217        if self.state.is_some() {
218            if let Err(error) = self.restore() {
219                tracing::warn!(%error, "failed to restore terminal modes");
220            }
221        }
222    }
223}
224
225impl Drop for TerminalModeRestoreGuard {
226    fn drop(&mut self) {
227        self.restore_silently();
228    }
229}
230
231pub(crate) struct TuiOptions<E> {
232    pub(crate) surface_preference: UiSurfacePreference,
233    pub(crate) inline_rows: u16,
234    pub(crate) show_logs: bool,
235    pub(crate) log_theme: Option<String>,
236    pub(crate) event_callback: Option<EventCallback<E>>,
237    pub(crate) focus_callback: Option<FocusChangeCallback>,
238    pub(crate) active_pty_sessions: Option<std::sync::Arc<std::sync::atomic::AtomicUsize>>,
239    pub(crate) input_activity_counter: Option<std::sync::Arc<std::sync::atomic::AtomicU64>>,
240    pub(crate) keyboard_protocol: crate::tui::config::KeyboardProtocolConfig,
241    pub(crate) fullscreen: FullscreenInteractionSettings,
242    pub(crate) workspace_root: Option<std::path::PathBuf>,
243    pub(crate) preview_callback: Option<PreviewCallback>,
244}
245
246pub(crate) async fn run_tui<S, F>(
247    mut commands: UnboundedReceiver<S::Command>,
248    events: UnboundedSender<S::Event>,
249    options: TuiOptions<S::Event>,
250    make_session: F,
251) -> Result<()>
252where
253    S: TuiSessionDriver,
254    F: FnOnce(u16) -> S,
255{
256    // Create a guard to mark TUI as initialized during the session
257    // This ensures the panic hook knows to restore terminal state
258    let _panic_guard = crate::tui::ui::tui::panic_hook::TuiPanicGuard::new();
259    crate::tui::frame_metrics::initialize_from_env();
260
261    let _signal_guard = SignalCleanupGuard::new()?;
262
263    let surface = TerminalSurface::detect(options.surface_preference, options.inline_rows)?;
264    set_log_theme_name(options.log_theme.clone());
265    let mut session = make_session(surface.rows());
266    session.attach_program_status_terminal();
267    session.set_preview_callback(options.preview_callback.clone());
268    session.set_show_logs(options.show_logs);
269    session.set_active_pty_sessions(options.active_pty_sessions);
270    session.set_workspace_root(options.workspace_root.clone());
271    session.set_fullscreen_active(surface.use_alternate());
272    session.set_fullscreen_interaction(options.fullscreen);
273    if options.show_logs {
274        let (log_tx, log_rx) = tokio::sync::mpsc::unbounded_channel();
275        session.set_log_receiver(log_rx);
276        register_tui_log_sender(log_tx);
277    } else {
278        clear_tui_log_sender();
279    }
280
281    let keyboard_flags = crate::tui::config::keyboard_protocol_to_flags(&options.keyboard_protocol);
282    let mut stderr = io::stderr();
283    let mut mode_restore_guard =
284        TerminalModeRestoreGuard::new(enable_terminal_modes(&mut stderr, &options.fullscreen)?);
285    mode_restore_guard.state_mut().save_cursor_position(&mut stderr);
286    if surface.use_alternate() {
287        mode_restore_guard.state_mut().enter_alternate_screen(&mut stderr)?;
288        // Record the surface so the canonical restore path can skip the
289        // full-screen clear: leaving the alternate buffer already restores
290        // the main screen (see panic_hook::restore_tui).
291        crate::tui::core_tui::panic_hook::state::mark_alternate_screen_active(true);
292    }
293    mode_restore_guard
294        .state_mut()
295        .push_keyboard_enhancement_flags(&mut stderr, keyboard_flags);
296
297    session.update_terminal_title();
298    super::session::terminal_title::apply_iterm2_profile_once();
299
300    let backend = CrosstermBackend::new(stderr);
301    let mut terminal = Terminal::new(backend).context("failed to initialize inline terminal")?;
302    prepare_terminal(&mut terminal)?;
303
304    // Create event listener and channels using the new async pattern
305    let (mut input_listener, event_channels) = EventListener::new();
306    let cancellation_token = CancellationToken::new();
307    let event_loop_token = cancellation_token.clone();
308    let event_channels_for_loop = event_channels.clone();
309    let rx_paused = event_channels.rx_paused.clone();
310    let last_input_elapsed_ms = event_channels.last_input_elapsed_ms.clone();
311    let session_start = event_channels.session_start;
312
313    // Ensure any capability or resize responses emitted during terminal setup are not treated as
314    // the user's first keystrokes.
315    drain_terminal_events();
316
317    // Clone the sender before moving event_channels_for_loop into tokio::spawn.
318    let event_tx_for_controller = event_channels_for_loop.tx.clone();
319
320    // Spawn the async event loop after the terminal is fully configured so the first keypress is
321    // captured immediately (avoids cooked-mode buffering before raw mode is enabled).
322    let event_loop_handle = tokio::spawn(async move {
323        spawn_event_loop(
324            event_channels_for_loop.tx.clone(),
325            event_loop_token,
326            rx_paused,
327            last_input_elapsed_ms,
328            session_start,
329        )
330        .await;
331    });
332
333    let mut event_stream = EventStreamController::new(
334        cancellation_token,
335        event_loop_handle,
336        event_tx_for_controller,
337        event_channels.rx_paused.clone(),
338        event_channels.last_input_elapsed_ms.clone(),
339        event_channels.session_start,
340    );
341
342    let drive_result = drive_terminal(
343        &mut terminal,
344        &mut session,
345        &mut commands,
346        &events,
347        &mut input_listener,
348        event_channels,
349        DriveRuntimeOptions {
350            event_callback: options.event_callback,
351            focus_callback: options.focus_callback,
352            use_alternate_screen: surface.use_alternate(),
353            input_activity_counter: options.input_activity_counter,
354            keyboard_flags,
355            fullscreen: options.fullscreen,
356            preview_callback: options.preview_callback,
357        },
358        &mut event_stream,
359    )
360    .await;
361
362    session.shutdown_program_status();
363
364    // Gracefully shutdown the event loop (may already be stopped by StopEventStream)
365    event_stream.shutdown().await;
366
367    // Drain any pending events before finalizing terminal and disabling modes
368    drain_terminal_events();
369
370    // When another party already restored the terminal (host backstop or panic
371    // hook), the main screen buffer is live: writes here (line clear, cursor
372    // show, clear) would land there instead of on the alternate screen.
373    let is_alternate = surface.use_alternate();
374    let finalize_result = if crate::tui::core_tui::panic_hook::is_restore_claimed() {
375        Ok(())
376    } else {
377        finalize_terminal(&mut terminal, is_alternate)
378    };
379
380    // Restore terminal modes (handles all modes including raw mode)
381    if let Err(error) = mode_restore_guard.restore() {
382        tracing::warn!(%error, "failed to restore terminal modes");
383    }
384
385    // Clear terminal title on exit.
386    session.clear_terminal_title();
387
388    drive_result?;
389    finalize_result?;
390
391    clear_tui_log_sender();
392    vtcode_commons::trace_flush::flush_trace_log();
393
394    Ok(())
395}
396
397#[cfg(test)]
398mod tests {
399    use super::*;
400    use std::sync::atomic::{AtomicBool, Ordering};
401
402    /// Sets a flag when dropped, so a test can observe that an aborted task
403    /// future was actually dropped (not merely detached).
404    struct DropFlag(std::sync::Arc<AtomicBool>);
405    impl Drop for DropFlag {
406        fn drop(&mut self) {
407            self.0.store(true, Ordering::SeqCst);
408        }
409    }
410
411    /// A wedged event-loop task must be aborted, not silently detached: the
412    /// prior implementation dropped the join handle on timeout, which left the
413    /// crossterm reader holding stdin past TUI teardown.
414    #[tokio::test]
415    async fn event_stream_shutdown_aborts_a_wedged_reader() {
416        let (_listener, channels) = EventListener::new();
417        let dropped = std::sync::Arc::new(AtomicBool::new(false));
418        let flag = dropped.clone();
419        let handle = tokio::spawn(async move {
420            let _guard = DropFlag(flag);
421            std::future::pending::<()>().await;
422        });
423
424        let mut controller = EventStreamController::new(
425            CancellationToken::new(),
426            handle,
427            channels.tx.clone(),
428            channels.rx_paused.clone(),
429            channels.last_input_elapsed_ms.clone(),
430            channels.session_start,
431        );
432
433        let started = std::time::Instant::now();
434        controller.shutdown().await;
435        assert!(started.elapsed() < Duration::from_secs(1), "shutdown must not block on a wedged reader");
436
437        tokio::time::timeout(Duration::from_secs(1), async {
438            while !dropped.load(Ordering::SeqCst) {
439                tokio::task::yield_now().await;
440            }
441        })
442        .await
443        .expect("shutdown must abort the wedged event-loop task");
444    }
445}