From fc44f61a10ac04688b2fedfa670f38a463ca25ac Mon Sep 17 00:00:00 2001 From: Can Celik Date: Sun, 13 Sep 2026 13:47:53 +0300 Subject: [PATCH] fix: disconnect stalled terminal observers (#4039) --- .../src/content/docs/cli-reference.mdx | 3 + src/platform/linux.rs | 2 +- src/platform/macos.rs | 2 +- src/platform/unix_common.rs | 89 ++++++++++++ src/server/client_transport.rs | 129 +++++++++++++++++- src/server/headless/render.rs | 10 ++ src/server/headless/tests/mod.rs | 47 +++++++ 7 files changed, 277 insertions(+), 5 deletions(-) diff --git a/docs/next/website/src/content/docs/cli-reference.mdx b/docs/next/website/src/content/docs/cli-reference.mdx index c25b64e6..ad9d7c5a 100644 --- a/docs/next/website/src/content/docs/cli-reference.mdx +++ b/docs/next/website/src/content/docs/cli-reference.mdx @@ -380,6 +380,9 @@ terminal, or agent target. It prints newline-delimited JSON `terminal.frame` records with base64-encoded ANSI bytes, then a `terminal.closed` record when the server closes the stream. Multiple observers can watch the same terminal without taking input, resize, scroll, or takeover authority. +On Linux and macOS, an observer is disconnected if a socket write makes no progress +for 30 seconds. A stalled stream can end without a final `terminal.closed` record; +the pane and other clients keep running. Quiet panes do not trigger this timeout. `terminal title clear` hands the outer terminal window title back to `ui.window_title`. ## Output waits diff --git a/src/platform/linux.rs b/src/platform/linux.rs index 29534dbe..23393653 100644 --- a/src/platform/linux.rs +++ b/src/platform/linux.rs @@ -17,7 +17,7 @@ pub(crate) use super::unix_common::{ create_remote_ssh_config_file, hostname, local_datetime, remote_bridge_endpoint_path, remote_private_temp_base, remote_reattach_argument, remote_reattach_program, remote_ssh_config_paths, set_default_plugin_pane_pwd, status_commands_supported, - wait_client_stream_readable, StatusCommandGuard, + wait_client_stream_readable, write_client_stream, ClientStreamReader, StatusCommandGuard, }; const WSL_MARKER_ENV_VARS: &[&str] = &["WSL_DISTRO_NAME", "WSL_INTEROP"]; diff --git a/src/platform/macos.rs b/src/platform/macos.rs index 8840c7ed..453d3fb3 100644 --- a/src/platform/macos.rs +++ b/src/platform/macos.rs @@ -17,7 +17,7 @@ pub(crate) use super::unix_common::{ create_remote_ssh_config_file, hostname, local_datetime, remote_bridge_endpoint_path, remote_private_temp_base, remote_reattach_argument, remote_reattach_program, remote_ssh_config_paths, set_default_plugin_pane_pwd, status_commands_supported, - wait_client_stream_readable, StatusCommandGuard, + wait_client_stream_readable, write_client_stream, ClientStreamReader, StatusCommandGuard, }; const PROC_PGRP_ONLY: u32 = 2; diff --git a/src/platform/unix_common.rs b/src/platform/unix_common.rs index 756abc17..ccb068f8 100644 --- a/src/platform/unix_common.rs +++ b/src/platform/unix_common.rs @@ -8,6 +8,95 @@ pub(crate) fn classify_child_exit(status: &portable_pty::ExitStatus) -> super::C } } +pub(crate) struct ClientStreamReader<'a>(pub(crate) &'a mut crate::ipc::LocalStream); + +impl std::io::Read for ClientStreamReader<'_> { + fn read(&mut self, data: &mut [u8]) -> std::io::Result { + use std::os::fd::AsRawFd as _; + + loop { + match self.0.read(data) { + Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => { + let crate::ipc::LocalStream::UdSocket(stream) = &*self.0; + let mut descriptor = libc::pollfd { + fd: stream.inner().as_raw_fd(), + events: libc::POLLIN, + revents: 0, + }; + // Sleep until input or shutdown, without polling quiet observers. + if unsafe { libc::poll(&mut descriptor, 1, -1) } < 0 { + let error = std::io::Error::last_os_error(); + if error.kind() != std::io::ErrorKind::Interrupted { + return Err(error); + } + } + } + result => return result, + } + } + } +} + +pub(crate) fn write_client_stream( + stream: &crate::ipc::LocalStream, + mut data: &[u8], +) -> std::io::Result<()> { + use std::io::{self, Write as _}; + use std::os::fd::AsRawFd as _; + use std::time::Instant; + + let crate::ipc::LocalStream::UdSocket(stream) = stream; + let mut socket = stream.inner(); + let Some(timeout) = socket.write_timeout()? else { + return socket.write_all(data); + }; + let timed_out = || { + // Dropping the writer clone alone would leave the reader blocked. + let _ = stream.inner().shutdown(std::net::Shutdown::Both); + io::Error::new( + io::ErrorKind::TimedOut, + "terminal observer stopped receiving output", + ) + }; + let mut progress = Instant::now(); + while !data.is_empty() { + match socket.write(data) { + Ok(0) => return Err(io::ErrorKind::WriteZero.into()), + Ok(written) => { + data = &data[written..]; + progress = Instant::now(); + continue; + } + Err(error) + if matches!( + error.kind(), + io::ErrorKind::WouldBlock | io::ErrorKind::Interrupted + ) => {} + Err(error) => return Err(error), + } + let remaining = timeout + .checked_sub(progress.elapsed()) + .ok_or_else(timed_out)?; + let mut descriptor = libc::pollfd { + fd: socket.as_raw_fd(), + events: libc::POLLOUT, + revents: 0, + }; + let wait_ms = remaining.as_millis().clamp(1, i32::MAX as u128) as i32; + let ready = unsafe { libc::poll(&mut descriptor, 1, wait_ms) }; + if ready == 0 { + return Err(timed_out()); + } + if ready < 0 { + let error = io::Error::last_os_error(); + if error.kind() != io::ErrorKind::Interrupted { + return Err(error); + } + } + } + Ok(()) +} + pub(crate) fn wait_client_stream_readable(stream: &crate::ipc::LocalStream) -> std::io::Result<()> { use std::os::fd::{AsFd as _, AsRawFd as _}; let crate::ipc::LocalStream::UdSocket(stream) = stream; diff --git a/src/server/client_transport.rs b/src/server/client_transport.rs index c57d0f86..da682c06 100644 --- a/src/server/client_transport.rs +++ b/src/server/client_transport.rs @@ -41,6 +41,9 @@ const MIN_CLIENT_ROWS: u16 = 1; /// and cleanup overhead. const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(4); +#[cfg(unix)] +const OBSERVER_WRITE_TIMEOUT: Duration = Duration::from_secs(30); + /// Maximum input payload size (bytes) for a single `ClientMessage::Input`. const MAX_INPUT_PAYLOAD: usize = 1024 * 1024; // 1 MB const MAX_CLIENT_SHELL_DIMENSION: u16 = 4096; @@ -940,7 +943,11 @@ fn client_writer_loop( } fn write_framed_bytes(stream: &mut LocalStream, data: &[u8]) -> bool { - if let Err(err) = stream.write_all(data) { + #[cfg(unix)] + let result = crate::platform::write_client_stream(stream, data); + #[cfg(windows)] + let result = stream.write_all(data); + if let Err(err) = result { debug!(err = %err, "client write failed, closing writer"); return false; } @@ -970,8 +977,14 @@ fn client_read_loop_with_endpoint_controls( endpoint_control_writer: Option<&ClientControlWriter>, ) -> io::Result<()> { while !should_quit.load(Ordering::Acquire) { - let msg: ClientMessage = match protocol::read_message(&mut stream, MAX_GRAPHICS_FRAME_SIZE) - { + #[cfg(unix)] + let message = protocol::read_message( + &mut crate::platform::ClientStreamReader(&mut stream), + MAX_GRAPHICS_FRAME_SIZE, + ); + #[cfg(windows)] + let message = protocol::read_message(&mut stream, MAX_GRAPHICS_FRAME_SIZE); + let msg: ClientMessage = match message { Ok(msg) => msg, Err(protocol::FramingError::UnexpectedEof) => { // Client disconnected. @@ -1027,6 +1040,13 @@ fn client_read_loop_with_endpoint_controls( } } ClientMessage::ObserveTerminal { target } => { + #[cfg(unix)] + { + stream.set_send_timeout(Some(OBSERVER_WRITE_TIMEOUT))?; + // macOS Unix sockets can block even with per-send MSG_DONTWAIT. + // ClientStreamReader preserves blocking reads on the shared socket. + stream.set_nonblocking(true)?; + } ServerEvent::ClientObserveTerminal { client_id, target } } ClientMessage::ControlTerminal { target, takeover } => { @@ -1641,6 +1661,109 @@ mod tests { )); } + #[cfg(unix)] + #[test] + fn observer_write_timeout_resets_when_sending_makes_progress() { + use std::io::Read as _; + + let (mut client, mut server, _path) = local_stream_pair("slow-observer"); + server + .set_send_timeout(Some(Duration::from_millis(100))) + .unwrap(); + server.set_nonblocking(true).unwrap(); + let worker = std::thread::spawn(move || { + assert!(write_framed_bytes(&mut server, &vec![b'x'; 1024 * 1024])); + }); + client + .set_recv_timeout(Some(Duration::from_secs(3))) + .unwrap(); + let mut received = 0; + let mut buffer = [0; 16 * 1024]; + while received < 1024 * 1024 { + let count = client.read(&mut buffer).unwrap(); + assert_ne!(count, 0, "observer disconnected while making progress"); + received += count; + std::thread::sleep(Duration::from_millis(5)); + } + worker.join().unwrap(); + } + + #[cfg(unix)] + #[test] + fn stalled_observer_timeout_releases_writer_and_reader() { + let (mut client, server, _path) = local_stream_pair("stalled-observer"); + let writer_stream = server.try_clone().expect("clone writer stream"); + let (writer, queue) = test_queue_writer(); + let (events, event_rx) = mpsc::channel(8); + let reader_events = events.clone(); + let (reader_done_tx, reader_done) = std::sync::mpsc::channel(); + let reader = std::thread::spawn(move || { + let result = client_read_loop( + server, + 14, + &reader_events, + &Arc::new(AtomicBool::new(false)), + ); + let _ = reader_done_tx.send(result); + }); + protocol::write_message( + &mut client, + &ClientMessage::ObserveTerminal { + target: "w1:p1".into(), + }, + ) + .expect("observe request"); + let mut event_rx = event_rx; + assert!(matches!( + event_rx.blocking_recv(), + Some(ServerEvent::ClientObserveTerminal { client_id: 14, .. }) + )); + let LocalStream::UdSocket(socket) = &writer_stream; + assert_eq!( + socket.inner().write_timeout().unwrap(), + Some(Duration::from_secs(30)) + ); + assert_eq!(socket.inner().read_timeout().unwrap(), None); + let mut resize = Vec::new(); + protocol::write_message( + &mut resize, + &ClientMessage::Resize { + cols: 80, + rows: 24, + cell_width_px: 0, + cell_height_px: 0, + pixel_mouse: false, + }, + ) + .unwrap(); + client.write_all(&resize[..2]).unwrap(); + std::thread::sleep(Duration::from_millis(20)); + client.write_all(&resize[2..]).unwrap(); + assert!(matches!( + event_rx.blocking_recv(), + Some(ServerEvent::ClientResize { client_id: 14, .. }) + )); + writer_stream + .set_send_timeout(Some(Duration::from_millis(200))) + .unwrap(); + let (writer_done_tx, writer_done) = std::sync::mpsc::channel(); + let worker = std::thread::spawn(move || { + client_writer_loop(writer_stream, 14, queue, events); + let _ = writer_done_tx.send(()); + }); + writer.render.try_send(vec![0; 4 * 1024 * 1024]).unwrap(); + writer_done + .recv_timeout(Duration::from_millis(350)) + .expect("writer timed out"); + reader_done + .recv_timeout(Duration::from_secs(2)) + .expect("reader released") + .unwrap(); + worker.join().unwrap(); + reader.join().unwrap(); + assert!(writer.control.send(vec![1]).is_err()); + } + #[test] fn clamp_terminal_size_zero_zero() { assert_eq!( diff --git a/src/server/headless/render.rs b/src/server/headless/render.rs index fc55417f..882a1615 100644 --- a/src/server/headless/render.rs +++ b/src/server/headless/render.rs @@ -462,6 +462,16 @@ impl HeadlessServer { let mut broken_clients: Vec = Vec::new(); for (client_id, (cols, rows), cell_size, _is_foreground, mode) in render_targets { + #[cfg(unix)] + if matches!(mode, ClientConnectionMode::TerminalObserve { .. }) + && self + .clients + .get(&client_id) + .is_some_and(|client| client.deferred_render() != DeferredRender::None) + { + // The writer-drained event schedules recovery, even if pane output stops. + continue; + } let area = Rect::new(0, 0, cols, rows); let shell_target = self.shell_target_for_client(client_id); let shell_tab_id = self.shell_tab_id_for_client(client_id); diff --git a/src/server/headless/tests/mod.rs b/src/server/headless/tests/mod.rs index 2c13dc8b..bb196be5 100644 --- a/src/server/headless/tests/mod.rs +++ b/src/server/headless/tests/mod.rs @@ -3509,6 +3509,53 @@ fn terminal_attach_disconnect_restores_client_shell_pane_size() { rt.shutdown_timeout(Duration::from_millis(100)); } +#[cfg(unix)] +#[test] +fn backpressured_observer_skips_runtime_access_and_recovers_without_new_output() { + with_terminal_session_test_server(|server, terminal_id, target, _| { + let (writer, control, frames) = test_client_writer(); + server.handle_server_event(ServerEvent::ClientConnected { + client_id: 7, + cols: 80, + rows: 24, + cell_width_px: 0, + cell_height_px: 0, + pixel_mouse: false, + writer, + }); + server.handle_server_event(ServerEvent::ClientObserveTerminal { + client_id: 7, + target, + }); + server.render_and_stream(); + server.app.terminal_runtimes.insert( + terminal_id.clone(), + crate::terminal::TerminalRuntime::test_with_screen_bytes(80, 24, b"LATEST"), + ); + server.render_and_stream(); + assert_eq!(server.clients[&7].deferred_render(), DeferredRender::Full); + + // An absent runtime makes any attempted rendering observable without timing a lock. + let runtime = server.app.terminal_runtimes.remove(&terminal_id).unwrap(); + server.render_and_stream(); + server.app.terminal_runtimes.insert(terminal_id, runtime); + assert!( + server.clients.contains_key(&7), + "backpressure must skip runtime access" + ); + assert!(control.try_recv().is_err()); + + let _ = frames.recv().expect("previously accepted frame"); + assert!(server.handle_server_event(ServerEvent::ClientWriterDrained { client_id: 7 })); + server.render_and_stream(); + let ServerMessage::Terminal(frame) = read_server_message(frames.recv().unwrap()) else { + panic!("terminal update"); + }; + assert!(String::from_utf8_lossy(&frame.bytes).contains("LATEST")); + assert_eq!(server.clients[&7].deferred_render(), DeferredRender::None); + }); +} + #[test] fn terminal_observe_allows_multiple_clients_without_attach_ownership() { with_terminal_session_test_server(|server, terminal_id, terminal_id_string, _| {