-
Notifications
You must be signed in to change notification settings - Fork 68
fix: drop attached clients that stop reading #436
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We鈥檒l occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This seems like generic functionality, so let's name it accordingly. |
||
| 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<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() }; | ||
|
|
@@ -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 = | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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<()> { | ||
|
|
||
There was a problem hiding this comment.
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.