vtcode_ui/tui/core_tui/runner/
mod.rs1use 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 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
116pub(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 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 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 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 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 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 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 drain_terminal_events();
316
317 let event_tx_for_controller = event_channels_for_loop.tx.clone();
319
320 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 event_stream.shutdown().await;
366
367 drain_terminal_events();
369
370 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 if let Err(error) = mode_restore_guard.restore() {
382 tracing::warn!(%error, "failed to restore terminal modes");
383 }
384
385 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 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 #[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}