From 8c73450962bedfd90a576061a6e0ca76f4ade941 Mon Sep 17 00:00:00 2001 From: warelik Date: Tue, 22 Sep 2026 06:29:18 +0000 Subject: [PATCH] fix: drop attached clients that stop reading A client that stays connected but stops draining its socket, e.g. an IDE terminal that was orphaned when its window lost the remote connection, used to park the shell->client thread in write() forever. Once the pty filled up the shell froze, and the session stayed attached: detach and attach -f both time out handing the disconnect to the blocked thread, so the session could only be recovered by finding and killing the stale attach process. Give writes to the client a 10 second timeout. When a write makes no progress for that long, treat the client as gone and shut its socket down, so the client->shell thread and the attach process see EOF, the session becomes disconnected, and the shell keeps running into the output spool. Stop the session-restore dump at the first failed chunk too, since each further chunk could otherwise block for another timeout. force_attach_stalled_client floods a session whose attach client never reads its stdout and checks that the session is released and the shell is still usable. It fails without this change. --- libshpool/src/daemon/shell.rs | 46 ++++++++++++++++++++++++++++++++ shpool/tests/attach.rs | 50 +++++++++++++++++++++++++++++++++++ 2 files changed, 96 insertions(+) diff --git a/libshpool/src/daemon/shell.rs b/libshpool/src/daemon/shell.rs index 56aa3d26..0b795df3 100644 --- a/libshpool/src/daemon/shell.rs +++ b/libshpool/src/daemon/shell.rs @@ -63,6 +63,14 @@ const SHELL_TO_CLIENT_POLL_MS: u16 = 50; // shell->client thread. const SHELL_TO_CLIENT_CTL_TIMEOUT: time::Duration = time::Duration::from_millis(300); +// How long a single write to an attached client may block before we +// decide the client has stopped reading and drop it. A client that stays +// connected but never drains its socket (an orphaned terminal emulator, +// a stalled ssh window) would otherwise park the shell->client thread in +// write() forever, freezing the shell once the pty fills up and keeping +// the session attached so that neither detach nor attach -f can take it. +const CLIENT_WRITE_TIMEOUT: time::Duration = time::Duration::from_secs(10); + /// Lifecycle state tracking when sessions were last connected/disconnected and /// by whom. The fields update in lockstep, so they live behind a single private /// lock that only this type's methods take, exactly once per update or read. @@ -233,6 +241,13 @@ fn shutdown_socket( } } +/// Reports whether a client write failed because it hit +/// CLIENT_WRITE_TIMEOUT. Depending on the platform, an expired +/// socket write timeout surfaces as either kind. +fn is_client_write_timeout(err: &io::Error) -> bool { + matches!(err.kind(), io::ErrorKind::WouldBlock | io::ErrorKind::TimedOut) +} + /// Messages to the shell->client thread to add or remove a client connection. pub enum ClientConnectionMsg { /// Accept a newly connected client @@ -432,6 +447,11 @@ impl SessionInner { trace!("client hangup: {:?}", e); false } + Err(e) if is_client_write_timeout(&e) => { + warn!("client stopped reading, dropping it: {:?}", e); + Self::drop_client(conn); + false + } Err(e) => { error!("unexpected IO error while writing heartbeat: {}", e); return Err(e).context("writing heartbeat")?; @@ -476,6 +496,10 @@ impl SessionInner { Err(e) if e.kind() == io::ErrorKind::BrokenPipe => { trace!("writing MaybeSwitch: client hangup: {:?}", e); } + Err(e) if is_client_write_timeout(&e) => { + warn!("writing MaybeSwitch: client stopped reading, dropping it: {:?}", e); + Self::drop_client(conn); + } Err(e) => { error!("unexpected IO error while writing heartbeat: {}", e); return Err(e).context("writing MaybeSwitch")?; @@ -522,6 +546,11 @@ impl SessionInner { if let Err(err) = chunk.write_to(&mut conn.sink) { warn!("err writing session-restore buf: {:?}", err); + // Each further chunk could block for another + // CLIENT_WRITE_TIMEOUT, so give up on the + // client right away. + Self::drop_client(conn); + break; } } if let Err(err) = conn.sink.flush() { @@ -624,6 +653,7 @@ impl SessionInner { chunk.write_to(&mut conn.sink).and_then(|_| conn.sink.flush()); if let Err(err) = write_result { info!("client_stream write err, assuming hangup: {:?}", err); + Self::drop_client(conn); reset_client_conn = true; } else { test_hooks::emit("daemon-wrote-s2c-chunk"); @@ -640,6 +670,16 @@ impl SessionInner { .spawn(move || log_if_error("error in shell->client", closure()))?) } + /// Gives up on a client we can no longer write to. The socket is shut + /// down rather than just dropped: the client may still be alive, just + /// not reading, and the client->shell thread and the attach process + /// only notice that the connection is over once they see EOF. + fn drop_client(conn: &ClientConnection) { + if let Err(e) = shutdown_socket(&conn.stream, net::Shutdown::Both) { + warn!("shutting down client stream: {:?}", e); + } + } + fn write_exit_chunk(mut sink: W, status: i32) { let status_buf: [u8; 4] = status.to_le_bytes(); let chunk = Chunk { kind: ChunkKind::ExitStatus, buf: status_buf.as_slice() }; @@ -677,6 +717,12 @@ impl SessionInner { None => return Err(anyhow!("no client stream to take for bidi streaming")), }; + // The timeout belongs to the socket, so it also covers the clones + // below that the shell->client thread writes through. + client_stream + .set_write_timeout(Some(CLIENT_WRITE_TIMEOUT)) + .context("setting client write timeout")?; + let mut client_to_shell_client_stream = client_stream.try_clone().context("creating client->shell client stream")?; let shell_to_client_client_stream = diff --git a/shpool/tests/attach.rs b/shpool/tests/attach.rs index e5a10402..53a114fd 100644 --- a/shpool/tests/attach.rs +++ b/shpool/tests/attach.rs @@ -923,6 +923,56 @@ fn force_attach() -> anyhow::Result<()> { Ok(()) } +#[test] +#[timeout(30000)] +fn force_attach_stalled_client() -> anyhow::Result<()> { + let mut daemon_proc = support::daemon::Proc::new( + "norc.toml", + DaemonArgs { listen_events: false, ..DaemonArgs::default() }, + ) + .context("starting daemon proc")?; + let disconnected_re = Regex::new("sh1.*disconnected").expect("valid regex"); + + let mut tty1 = daemon_proc.attach("sh1", Default::default()).context("attaching from tty1")?; + let mut line_matcher1 = tty1.line_matcher()?; + tty1.run_cmd("export MYVAR='set_from_tty1'")?; + tty1.run_cmd("echo $MYVAR")?; + line_matcher1.scan_until_re("set_from_tty1$")?; + + // Flood the session and never read tty1's stdout again, like a + // terminal that is still open but no longer drained (an orphaned + // IDE terminal or a stalled ssh window). tty1 blocks writing its + // stdout and stops reading from the daemon. + tty1.run_cmd("yes")?; + + // The daemon has to give up on the stalled client on its own + // rather than letting it hold the session (and the shell's + // output) forever. That takes the daemon's client write timeout, + // and the exponential backoff in wait_until_list_matches polls too + // sparsely around that point, so poll at a fixed interval instead. + let deadline = time::Instant::now() + time::Duration::from_secs(20); + loop { + let list_out = daemon_proc.list()?; + if disconnected_re.is_match(&String::from_utf8_lossy(&list_out.stdout[..])) { + break; + } + if time::Instant::now() > deadline { + return Err(anyhow!("the stalled client was never dropped")); + } + thread::sleep(time::Duration::from_millis(500)); + } + + let mut tty2 = daemon_proc + .attach("sh1", AttachArgs { force: true, ..Default::default() }) + .context("attaching from tty2")?; + let mut line_matcher2 = tty2.line_matcher()?; + tty2.run_raw(vec![3])?; // ^C the flood + tty2.run_cmd("echo $MYVAR")?; + line_matcher2.scan_until_re("set_from_tty1$")?; + + Ok(()) +} + #[test] #[timeout(30000)] fn busy() -> anyhow::Result<()> {