fix: disconnect stalled terminal observers (#4039)

This commit is contained in:
Can Celik
2026-09-13 13:47:53 +03:00
committed by GitHub
parent bafbc0949d
commit fc44f61a10
7 changed files with 277 additions and 5 deletions
@@ -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
+1 -1
View File
@@ -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"];
+1 -1
View File
@@ -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;
+89
View File
@@ -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<usize> {
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;
+126 -3
View File
@@ -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!(
+10
View File
@@ -462,6 +462,16 @@ impl HeadlessServer {
let mut broken_clients: Vec<u64> = 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);
+47
View File
@@ -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, _| {