Recover from sustained terminal backpressure
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / macos-client (macos-aarch64, macos-14) (push) Canceled after 0s
ci / macos-client (macos-x86_64, macos-13) (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest, , , ) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14, , , ) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13, , , ) (push) Canceled after 0s
ci / package-release (windows-aarch64, windows-latest, aarch64, windows, aarch64-pc-windows-msvc) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest, , , ) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
This commit is contained in:
+64
-19
@@ -6570,6 +6570,8 @@ async fn run_terminal(
|
||||
predict_mode,
|
||||
);
|
||||
let mut frame_renderer = TerminalFrameRenderer::new()?;
|
||||
let mut render_resync_needed = false;
|
||||
let mut overflow_closed_frame: Option<Frame> = None;
|
||||
// Non-destructive disconnect status line. Off in forward-only mode (no TTY to
|
||||
// draw on) and respecting the client config (default on, env override).
|
||||
let mut disconnect_status = DisconnectStatus::new(resolve_disconnect_status() && !forward_only);
|
||||
@@ -6676,14 +6678,16 @@ async fn run_terminal(
|
||||
let mut detach_requested = false;
|
||||
loop {
|
||||
accepted_output_seq = accepted_output_seq.max(cred.last_rendered_seq);
|
||||
predictor.set_output_backpressured(frame_renderer.has_pending());
|
||||
predictor.set_output_backpressured(frame_renderer.has_pending() || render_resync_needed);
|
||||
if detach_requested {
|
||||
break;
|
||||
}
|
||||
tokio::select! {
|
||||
rendered = frame_renderer.recv(), if frame_renderer.has_pending() => {
|
||||
let frame = rendered?;
|
||||
predictor.set_output_backpressured(frame_renderer.has_pending());
|
||||
predictor.set_output_backpressured(
|
||||
frame_renderer.has_pending() || render_resync_needed,
|
||||
);
|
||||
cred.last_rendered_seq = cred.last_rendered_seq.max(frame.output_seq);
|
||||
note_snapshot_rendered(
|
||||
&frame,
|
||||
@@ -6696,6 +6700,16 @@ async fn run_terminal(
|
||||
if frame.closed {
|
||||
return Ok(());
|
||||
}
|
||||
if !frame_renderer.has_pending()
|
||||
&& let Some(frame) = overflow_closed_frame.take()
|
||||
{
|
||||
render_resync_needed = false;
|
||||
predictor.observe_output(&frame.bytes);
|
||||
anyhow::ensure!(
|
||||
frame_renderer.enqueue(frame)?.is_none(),
|
||||
"terminal renderer remained full after draining"
|
||||
);
|
||||
}
|
||||
}
|
||||
stdin_msg = stdin_rx.recv() => {
|
||||
match stdin_msg {
|
||||
@@ -7098,7 +7112,8 @@ async fn run_terminal(
|
||||
maybe_send_resize(&socket, addr, &cred, &mut send_seq, &mut last_size).await?;
|
||||
}
|
||||
_ = frame_gap_tick.tick() => {
|
||||
if !frame_renderer.has_pending()
|
||||
if !render_resync_needed
|
||||
&& !frame_renderer.has_pending()
|
||||
&& frame_buffer.resync_due()
|
||||
&& let Some(frame) = reconnect(
|
||||
&socket,
|
||||
@@ -7219,8 +7234,30 @@ async fn run_terminal(
|
||||
predictor.observe_output(&frame.bytes);
|
||||
last_terminal_frame_at = Instant::now();
|
||||
wake_repaint_retry_until = None;
|
||||
frame_renderer.enqueue(frame)?;
|
||||
predictor.set_output_backpressured(true);
|
||||
if render_resync_needed {
|
||||
if frame.closed {
|
||||
overflow_closed_frame = Some(frame);
|
||||
}
|
||||
} else if let Some(frame) = frame_renderer.enqueue(frame)? {
|
||||
dosh::trace::event(
|
||||
"client.render_queue_overflow",
|
||||
&[
|
||||
("output_seq", frame.output_seq.to_string()),
|
||||
("bytes", frame.bytes.len().to_string()),
|
||||
],
|
||||
);
|
||||
dosh::trace::health_event(
|
||||
"client.render_queue_overflow",
|
||||
&[("output_seq", frame.output_seq.to_string())],
|
||||
);
|
||||
render_resync_needed = true;
|
||||
if frame.closed {
|
||||
overflow_closed_frame = Some(frame);
|
||||
}
|
||||
}
|
||||
predictor.set_output_backpressured(
|
||||
frame_renderer.has_pending() || render_resync_needed,
|
||||
);
|
||||
flush_startup_input_if_ready(
|
||||
&socket,
|
||||
addr,
|
||||
@@ -7945,7 +7982,8 @@ async fn run_terminal(
|
||||
let mut repainted_this_tick = false;
|
||||
let stale = last_packet_at.elapsed();
|
||||
if !frame_renderer.has_pending()
|
||||
&& (status_restore_pending
|
||||
&& (render_resync_needed
|
||||
|| status_restore_pending
|
||||
|| stale >= Duration::from_secs(reconnect_timeout_secs.max(1)))
|
||||
{
|
||||
if let Some(frame) = reconnect(
|
||||
@@ -7974,6 +8012,7 @@ async fn run_terminal(
|
||||
last_idle_repaint_attempt_at = Instant::now();
|
||||
wake_repaint_retry_until = None;
|
||||
repainted_this_tick = true;
|
||||
render_resync_needed = false;
|
||||
}
|
||||
last_packet_at = Instant::now();
|
||||
flush_pending_user_input(
|
||||
@@ -8004,7 +8043,7 @@ async fn run_terminal(
|
||||
// a latency spike (or recovery) flips speculation on/off promptly
|
||||
// without waiting for the next keystroke to drive `redraw`.
|
||||
if !forward_only {
|
||||
if !frame_renderer.has_pending() {
|
||||
if !frame_renderer.has_pending() && !render_resync_needed {
|
||||
predictor.refresh_policy()?;
|
||||
}
|
||||
flush_startup_input_if_ready(
|
||||
@@ -8018,7 +8057,7 @@ async fn run_terminal(
|
||||
)
|
||||
.await?;
|
||||
let now = Instant::now();
|
||||
if !repainted_this_tick && !frame_renderer.has_pending() && should_repaint_idle_terminal(
|
||||
if !repainted_this_tick && !render_resync_needed && !frame_renderer.has_pending() && should_repaint_idle_terminal(
|
||||
predictor.alternate_screen,
|
||||
last_terminal_frame_at,
|
||||
last_idle_repaint_attempt_at,
|
||||
@@ -8076,7 +8115,10 @@ async fn run_terminal(
|
||||
// on how long the link has been silent (recomputed after any
|
||||
// reconnect attempt above may have reset `last_packet_at`).
|
||||
if !forward_only {
|
||||
let action = if predictor.alternate_screen || frame_renderer.has_pending() {
|
||||
let action = if predictor.alternate_screen
|
||||
|| frame_renderer.has_pending()
|
||||
|| render_resync_needed
|
||||
{
|
||||
disconnect_status.on_suppressed()
|
||||
} else {
|
||||
disconnect_status.on_tick(last_packet_at.elapsed())
|
||||
@@ -10576,19 +10618,22 @@ impl TerminalFrameRenderer {
|
||||
})
|
||||
}
|
||||
|
||||
fn enqueue(&mut self, frame: Frame) -> Result<()> {
|
||||
/// Returns the frame unchanged when the bounded queue is full. The caller
|
||||
/// can continue servicing transport traffic and request a snapshot once the
|
||||
/// renderer drains instead of turning local display backpressure into a
|
||||
/// disconnected session.
|
||||
fn enqueue(&mut self, frame: Frame) -> Result<Option<Frame>> {
|
||||
let Some(jobs) = self.jobs.as_ref() else {
|
||||
return Err(anyhow!("terminal renderer is closed"));
|
||||
};
|
||||
jobs.try_send(frame).map_err(|err| match err {
|
||||
mpsc::error::TrySendError::Full(_) => anyhow!(
|
||||
"terminal renderer exceeded its bounded {}-frame queue",
|
||||
TERMINAL_RENDER_QUEUE_CAPACITY
|
||||
),
|
||||
mpsc::error::TrySendError::Closed(_) => anyhow!("terminal renderer stopped"),
|
||||
})?;
|
||||
self.pending += 1;
|
||||
Ok(())
|
||||
match jobs.try_send(frame) {
|
||||
Ok(()) => {
|
||||
self.pending += 1;
|
||||
Ok(None)
|
||||
}
|
||||
Err(mpsc::error::TrySendError::Full(frame)) => Ok(Some(frame)),
|
||||
Err(mpsc::error::TrySendError::Closed(_)) => Err(anyhow!("terminal renderer stopped")),
|
||||
}
|
||||
}
|
||||
|
||||
fn has_pending(&self) -> bool {
|
||||
|
||||
@@ -51,6 +51,7 @@ enum ServerObservation {
|
||||
BulkSent,
|
||||
Input(Vec<u8>),
|
||||
Reconnected,
|
||||
RenderResynced,
|
||||
Resize(u16, u16),
|
||||
}
|
||||
|
||||
@@ -67,7 +68,7 @@ fn native_client_terminal_round_trip_is_platform_complete() {
|
||||
.set_read_timeout(Some(Duration::from_secs(8)))
|
||||
.unwrap();
|
||||
let port = socket.local_addr().unwrap().port();
|
||||
write_client_fixture(&home, &cache, port);
|
||||
write_client_fixture(&home, &cache, port, 5);
|
||||
|
||||
let (observation_tx, observation_rx) = mpsc::channel();
|
||||
let server = thread::spawn(move || run_fake_terminal_server(socket, observation_tx));
|
||||
@@ -176,7 +177,7 @@ fn terminal_output_backpressure_does_not_block_input_transport() {
|
||||
.set_read_timeout(Some(Duration::from_secs(8)))
|
||||
.unwrap();
|
||||
let port = socket.local_addr().unwrap().port();
|
||||
write_client_fixture(&home, &cache, port);
|
||||
write_client_fixture(&home, &cache, port, 5);
|
||||
let (observation_tx, observation_rx) = mpsc::channel();
|
||||
let server = thread::spawn(move || run_backpressure_server(socket, observation_tx, PROBE));
|
||||
|
||||
@@ -240,7 +241,7 @@ fn idle_reconnect_restores_snapshot_and_orders_reordered_frames() {
|
||||
.set_read_timeout(Some(Duration::from_secs(8)))
|
||||
.unwrap();
|
||||
let port = socket.local_addr().unwrap().port();
|
||||
write_client_fixture(&home, &cache, port);
|
||||
write_client_fixture(&home, &cache, port, 1);
|
||||
let (observation_tx, observation_rx) = mpsc::channel();
|
||||
let server = thread::spawn(move || run_reconnect_server(socket, observation_tx, PROBE));
|
||||
|
||||
@@ -299,13 +300,92 @@ fn idle_reconnect_restores_snapshot_and_orders_reordered_frames() {
|
||||
assert!(status.success(), "Dosh client exited with {status:?}");
|
||||
}
|
||||
|
||||
fn write_client_fixture(home: &Path, cache: &Path, port: u16) {
|
||||
#[test]
|
||||
fn renderer_overflow_resyncs_without_dropping_the_session() {
|
||||
const PROBE: &[u8] = b"DOSH_OVERFLOW_INPUT\r";
|
||||
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let home = dir.path().join("home");
|
||||
let cache = dir.path().join("credentials");
|
||||
fs::create_dir_all(home.join(".config/dosh")).unwrap();
|
||||
fs::create_dir_all(&cache).unwrap();
|
||||
|
||||
let socket = UdpSocket::bind("127.0.0.1:0").unwrap();
|
||||
socket
|
||||
.set_read_timeout(Some(Duration::from_secs(12)))
|
||||
.unwrap();
|
||||
let port = socket.local_addr().unwrap().port();
|
||||
write_client_fixture(&home, &cache, port, 30);
|
||||
let (observation_tx, observation_rx) = mpsc::channel();
|
||||
let server = thread::spawn(move || run_overflow_server(socket, observation_tx, PROBE));
|
||||
|
||||
let pty = NativePtySystem::default();
|
||||
let pair = pty
|
||||
.openpty(PtySize {
|
||||
rows: 24,
|
||||
cols: 80,
|
||||
pixel_width: 0,
|
||||
pixel_height: 0,
|
||||
})
|
||||
.unwrap();
|
||||
let mut reader = pair.master.try_clone_reader().unwrap();
|
||||
let mut writer = pair.master.take_writer().unwrap();
|
||||
let mut command = client_command(dir.path(), port);
|
||||
command.env("HOME", home.to_string_lossy().to_string());
|
||||
command.env("USERPROFILE", home.to_string_lossy().to_string());
|
||||
command.env("APPDATA", home.to_string_lossy().to_string());
|
||||
command.env("LOCALAPPDATA", home.to_string_lossy().to_string());
|
||||
command.env("TERM", "xterm-256color");
|
||||
let mut child = pair.slave.spawn_command(command).unwrap();
|
||||
drop(pair.slave);
|
||||
|
||||
wait_for_observation(&observation_rx, Duration::from_secs(5), |observation| {
|
||||
matches!(observation, ServerObservation::BulkSent)
|
||||
});
|
||||
writer.write_all(PROBE).unwrap();
|
||||
writer.flush().unwrap();
|
||||
wait_for_input(&observation_rx, PROBE, Duration::from_secs(2));
|
||||
|
||||
let output = Arc::new(Mutex::new(Vec::new()));
|
||||
let reader_output = Arc::clone(&output);
|
||||
let reader_thread = thread::spawn(move || {
|
||||
let mut buf = [0u8; 16 * 1024];
|
||||
loop {
|
||||
match reader.read(&mut buf) {
|
||||
Ok(0) | Err(_) => break,
|
||||
Ok(n) => reader_output.lock().unwrap().extend_from_slice(&buf[..n]),
|
||||
}
|
||||
}
|
||||
});
|
||||
wait_for_observation(&observation_rx, Duration::from_secs(5), |observation| {
|
||||
matches!(observation, ServerObservation::RenderResynced)
|
||||
});
|
||||
wait_for_output(
|
||||
&output,
|
||||
b"DOSH_OVERFLOW_RESYNC_DOSH_OVERFLOW_DONE",
|
||||
Duration::from_secs(5),
|
||||
);
|
||||
|
||||
let status = child.wait().unwrap();
|
||||
drop(writer);
|
||||
drop(pair.master);
|
||||
reader_thread.join().unwrap();
|
||||
server.join().unwrap();
|
||||
assert!(status.success(), "Dosh client exited with {status:?}");
|
||||
}
|
||||
|
||||
fn write_client_fixture(
|
||||
home: &Path,
|
||||
cache: &Path,
|
||||
port: u16,
|
||||
reconnect_timeout_secs: u64,
|
||||
) {
|
||||
let config = ClientConfig {
|
||||
server: "local".to_string(),
|
||||
dosh_host: Some("127.0.0.1".to_string()),
|
||||
dosh_port: port,
|
||||
cache_attach_tickets: false,
|
||||
reconnect_timeout_secs: 1,
|
||||
reconnect_timeout_secs,
|
||||
credential_cache: cache.to_string_lossy().to_string(),
|
||||
auth_preference: "ssh".to_string(),
|
||||
predict: false,
|
||||
@@ -638,6 +718,101 @@ fn run_reconnect_server(
|
||||
}
|
||||
}
|
||||
|
||||
fn run_overflow_server(
|
||||
socket: UdpSocket,
|
||||
observations: mpsc::Sender<ServerObservation>,
|
||||
expected_input: &[u8],
|
||||
) {
|
||||
const BULK_FRAMES: u64 = 640;
|
||||
|
||||
let mut resume_count = 0u64;
|
||||
let mut bulk_sent = false;
|
||||
let mut resynced = false;
|
||||
let mut closed_sent = false;
|
||||
let mut received_input = Vec::new();
|
||||
let mut buf = [0u8; 65535];
|
||||
loop {
|
||||
let (n, source) = socket.recv_from(&mut buf).unwrap();
|
||||
let packet = protocol::decode(&buf[..n]).unwrap();
|
||||
match packet.header.kind {
|
||||
PacketKind::ResumeRequest => {
|
||||
resume_count += 1;
|
||||
if resume_count == 1 {
|
||||
send_frame(
|
||||
&socket,
|
||||
source,
|
||||
PacketKind::ResumeOk,
|
||||
1,
|
||||
10,
|
||||
b"DOSH_OVERFLOW_READY",
|
||||
true,
|
||||
false,
|
||||
);
|
||||
} else {
|
||||
send_frame(
|
||||
&socket,
|
||||
source,
|
||||
PacketKind::ResumeOk,
|
||||
700,
|
||||
700,
|
||||
b"DOSH_OVERFLOW_RESYNC_",
|
||||
true,
|
||||
false,
|
||||
);
|
||||
resynced = true;
|
||||
observations
|
||||
.send(ServerObservation::RenderResynced)
|
||||
.unwrap();
|
||||
}
|
||||
}
|
||||
PacketKind::Ack if !bulk_sent => {
|
||||
let bulk = [b'x'; 1024];
|
||||
for index in 0..BULK_FRAMES {
|
||||
send_frame(
|
||||
&socket,
|
||||
source,
|
||||
PacketKind::Frame,
|
||||
2 + index,
|
||||
11 + index,
|
||||
&bulk,
|
||||
false,
|
||||
false,
|
||||
);
|
||||
thread::sleep(Duration::from_millis(1));
|
||||
}
|
||||
bulk_sent = true;
|
||||
observations.send(ServerObservation::BulkSent).unwrap();
|
||||
}
|
||||
PacketKind::Ack if resynced && !closed_sent => {
|
||||
send_frame(
|
||||
&socket,
|
||||
source,
|
||||
PacketKind::Frame,
|
||||
701,
|
||||
701,
|
||||
b"DOSH_OVERFLOW_DONE",
|
||||
false,
|
||||
true,
|
||||
);
|
||||
closed_sent = true;
|
||||
}
|
||||
PacketKind::Input => {
|
||||
let plain =
|
||||
protocol::decrypt_body(&packet, &SESSION_KEY, CLIENT_TO_SERVER).unwrap();
|
||||
let input: Input = protocol::from_body(&plain).unwrap();
|
||||
received_input.extend_from_slice(&input.bytes);
|
||||
observations
|
||||
.send(ServerObservation::Input(input.bytes))
|
||||
.unwrap();
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
if closed_sent && contains(&received_input, expected_input) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
fn send_frame(
|
||||
socket: &UdpSocket,
|
||||
|
||||
Reference in New Issue
Block a user