diff --git a/src/bin/dosh-client.rs b/src/bin/dosh-client.rs index 2fc4cb3..268ca4d 100644 --- a/src/bin/dosh-client.rs +++ b/src/bin/dosh-client.rs @@ -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 = 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> { 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 { diff --git a/tests/client_terminal_runtime.rs b/tests/client_terminal_runtime.rs index c223faa..ece1879 100644 --- a/tests/client_terminal_runtime.rs +++ b/tests/client_terminal_runtime.rs @@ -51,6 +51,7 @@ enum ServerObservation { BulkSent, Input(Vec), 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, + 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,