Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 46 additions & 0 deletions libshpool/src/daemon/shell.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

10s seems like a solid default, but let's make this optionally configurable by adding an entry to config.rs. That will also allow the test to run in a more reasonable timeframe.


/// 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.
Expand Down Expand Up @@ -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 {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This seems like generic functionality, so let's name it accordingly. is_socket_write_timeout would describe what it does better.

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
Expand Down Expand Up @@ -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")?;
Expand Down Expand Up @@ -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")?;
Expand Down Expand Up @@ -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() {
Expand Down Expand Up @@ -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");
Expand All @@ -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<W: io::Write>(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() };
Expand Down Expand Up @@ -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 =
Expand Down
50 changes: 50 additions & 0 deletions shpool/tests/attach.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);

@ethanpailes ethanpailes Sep 22, 2026 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Having the test block for 10s is not great. The test suite already takes a while to run because of all the process juggling, and this guy will immediately become a long pole. To address this, and also make this feature more configurable, let's make this write timeout an option in the config file, and then set it to something short like 1s in this test.

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: can you give a useful name to this magic number with a constant?

tty2.run_cmd("echo $MYVAR")?;
line_matcher2.scan_until_re("set_from_tty1$")?;

Ok(())
}

#[test]
#[timeout(30000)]
fn busy() -> anyhow::Result<()> {
Expand Down
Loading