From 97cf165527c02982d18ba17f224f1f1fba1c2c75 Mon Sep 17 00:00:00 2001
From: DuProcess <273172371+DuProcess@users.noreply.github.com>
Date: Fri, 17 Jul 2026 21:06:26 -0400
Subject: [PATCH] Recover from sustained terminal backpressure
---
src/bin/dosh-client.rs | 83 ++++++++++----
tests/client_terminal_runtime.rs | 185 ++++++++++++++++++++++++++++++-
2 files changed, 244 insertions(+), 24 deletions(-)
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