From eb0c1db837e7c6ab43468409e1676b35b8c301a1 Mon Sep 17 00:00:00 2001 From: DuProcess <273172371+DuProcess@users.noreply.github.com> Date: Fri, 17 Jul 2026 22:01:02 -0400 Subject: [PATCH] Frame terminal reports across input reads --- src/bin/dosh-client.rs | 762 ++++++++++++++++++++----------- tests/client_terminal_runtime.rs | 9 +- 2 files changed, 509 insertions(+), 262 deletions(-) diff --git a/src/bin/dosh-client.rs b/src/bin/dosh-client.rs index a655e47..31f4d5c 100644 --- a/src/bin/dosh-client.rs +++ b/src/bin/dosh-client.rs @@ -72,11 +72,11 @@ const STREAM_RETIRED_TOMBSTONES: usize = 16 * 1024; const STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS: u32 = 8; const MAX_PENDING_USER_INPUT_BYTES: usize = 1024 * 1024; const STDIN_QUEUE_CAPACITY: usize = 64; +const MAX_TERMINAL_REPORT_BYTES: usize = 256; const STARTUP_INPUT_HOLD: Duration = Duration::from_millis(750); const POST_SUBMIT_ALL_INPUT_HOLD: Duration = Duration::from_millis(120); const STALE_TERMINAL_INPUT_AFTER: Duration = Duration::from_secs(2); -const POST_RECONNECT_STALE_INPUT_GRACE: Duration = Duration::from_secs(5); -const FOCUS_REPAINT_COOLDOWN: Duration = Duration::from_secs(1); +const TERMINAL_REPORT_HOLD: Duration = Duration::from_millis(12); const LOCAL_SLEEP_REPAINT_AFTER: Duration = Duration::from_secs(5); const LOCAL_SLEEP_REPAINT_RETRY_AFTER: Duration = Duration::from_secs(1); const LOCAL_SLEEP_REPAINT_RETRY_WINDOW: Duration = Duration::from_secs(10); @@ -6498,7 +6498,7 @@ where } } -fn terminal_input_channel() -> (mpsc::Sender>, mpsc::Receiver>) { +fn terminal_input_channel() -> (mpsc::Sender, mpsc::Receiver) { mpsc::channel(STDIN_QUEUE_CAPACITY) } @@ -6542,11 +6542,13 @@ async fn run_terminal( let mut stream_retransmit_tick = tokio::time::interval(dosh::transport::ADAPTIVE_RETRANSMIT_MIN); let mut frame_gap_tick = tokio::time::interval(Duration::from_millis(250)); + let mut terminal_report_tick = tokio::time::interval(TERMINAL_REPORT_HOLD); for interval in [ &mut status_tick, &mut resize_tick, &mut stream_retransmit_tick, &mut frame_gap_tick, + &mut terminal_report_tick, ] { interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); } @@ -6612,8 +6614,9 @@ async fn run_terminal( let mut pending_user_input_bytes = 0usize; let mut startup_input_hold_until: Option = None; let mut startup_gate_mode = StartupGateMode::HoldControl; - let mut stale_terminal_input_suppress_until: Option = None; - let mut last_focus_repaint_at = Instant::now() - FOCUS_REPAINT_COOLDOWN; + let mut stale_terminal_input_suppress_through: Option = None; + let mut terminal_input_filter = TerminalInputFilter::default(); + let mut deferred_filtered_input: Option = None; let mut last_idle_repaint_attempt_at = Instant::now() - LOCAL_SLEEP_REPAINT_RETRY_AFTER; let mut last_status_tick_at = Instant::now(); let mut wake_repaint_retry_until: Option = None; @@ -6644,7 +6647,8 @@ async fn run_terminal( startup_gate_mode = StartupGateMode::HoldAll; } - let (stdin_tx, mut stdin_rx) = terminal_input_channel(); + let (stdin_tx, mut stdin_rx) = terminal_input_channel::(); + let filtered_input_tx = stdin_tx.clone(); let _stdin_keepalive = if forward_only { Some(stdin_tx) } else { @@ -6657,7 +6661,13 @@ async fn run_terminal( match stdin.read(&mut buf) { Ok(0) => break, Ok(n) => { - if stdin_tx.blocking_send(buf[..n].to_vec()).is_err() { + if stdin_tx + .blocking_send(TerminalInputEvent::Raw { + bytes: buf[..n].to_vec(), + read_at: Instant::now(), + }) + .is_err() + { break; } } @@ -6707,13 +6717,62 @@ async fn run_terminal( ); } } + _ = terminal_report_tick.tick(), if terminal_input_filter.has_pending() || deferred_filtered_input.is_some() => { + let batch = deferred_filtered_input + .take() + .or_else(|| terminal_input_filter.flush_due(Instant::now())); + if let Some(batch) = batch + && let Err(err) = filtered_input_tx.try_send(TerminalInputEvent::Filtered(batch)) + && let TerminalInputEvent::Filtered(batch) = err.into_inner() + { + deferred_filtered_input = Some(batch); + } + } stdin_msg = stdin_rx.recv() => { match stdin_msg { - Some(mut bytes) => { - dosh::trace::event( - "client.stdin", - &[("bytes", dosh::trace::bytes_summary(&bytes))], - ); + Some(event) => { + let batch = match event { + TerminalInputEvent::Raw { bytes, read_at } => { + dosh::trace::event( + "client.stdin", + &[("bytes", dosh::trace::bytes_summary(&bytes))], + ); + let suppress_mouse = should_strip_stale_terminal_reports( + last_packet_at.elapsed(), + read_at, + stale_terminal_input_suppress_through, + ); + if stale_terminal_input_suppress_through + .is_some_and(|cutoff| read_at > cutoff) + { + stale_terminal_input_suppress_through = None; + } + terminal_input_filter.push( + bytes, + TerminalReportPolicy { + mouse_tracking: predictor.mouse_tracking.enabled(), + focus_tracking: predictor.focus_tracking, + suppress_mouse, + }, + read_at, + ) + } + TerminalInputEvent::Filtered(batch) => batch, + }; + let TerminalInputBatch { + mut bytes, + saw_focus_in, + stripped, + } = batch; + if stripped { + dosh::trace::event( + "client.terminal_reports_stripped", + &[("remaining", dosh::trace::bytes_summary(&bytes))], + ); + } + if bytes.is_empty() && !saw_focus_in { + continue; + } let (prefix_len, escape_found) = input_prefix_before_escape(&bytes, escape_key.as_deref()); if escape_found { @@ -6725,7 +6784,6 @@ async fn run_terminal( } } let input_status_tick_gap = last_status_tick_at.elapsed(); - let mut refreshed_before_input = false; if should_reconnect_before_input_for_local_sleep( forward_only, input_status_tick_gap, @@ -6747,9 +6805,6 @@ async fn run_terminal( )], ); last_status_tick_at = reconnect_started_at; - arm_stale_terminal_input_suppression( - &mut stale_terminal_input_suppress_until, - ); if let Some(frame) = reconnect( &socket, &mut cred, @@ -6768,64 +6823,8 @@ async fn run_terminal( ], ); refresh_live_addr(&mut addr, &cred)?; - render_frame(&frame)?; - note_snapshot_rendered( - &frame, - &mut disconnect_status, - &mut status_restore_pending, - ); - predictor.observe_output(&frame.bytes); - last_packet_at = Instant::now(); - last_focus_repaint_at = Instant::now(); - wake_repaint_retry_until = None; - refreshed_before_input = true; - flush_pending_user_input( - &socket, - addr, - &cred, - &mut send_seq, - &mut predictor, - &mut pending_user_input, - &mut pending_user_input_bytes, - ) - .await?; - } - } - let saw_focus_in = input_contains_focus_in(&bytes); - if !forward_only - && !refreshed_before_input - && !frame_renderer.has_pending() - && saw_focus_in - && last_focus_repaint_at.elapsed() >= FOCUS_REPAINT_COOLDOWN - { - dosh::trace::event( - "client.focus_reconnect_start", - &[("silent_ms", last_packet_at.elapsed().as_millis().to_string())], - ); - last_focus_repaint_at = Instant::now(); - wake_repaint_retry_until = - Some(last_focus_repaint_at + LOCAL_SLEEP_REPAINT_RETRY_WINDOW); - trace_wake_repaint_retry_arm("focus", last_packet_at.elapsed()); - if let Some(frame) = reconnect( - &socket, - &mut cred, - &mut send_seq, - last_size, - &mut frame_buffer, - &mut predictor, - ) - .await? - { - dosh::trace::event( - "client.focus_reconnect_ok", - &[ - ("output_seq", frame.output_seq.to_string()), - ("bytes", frame.bytes.len().to_string()), - ], - ); - refresh_live_addr(&mut addr, &cred)?; arm_stale_terminal_input_suppression( - &mut stale_terminal_input_suppress_until, + &mut stale_terminal_input_suppress_through, ); render_frame(&frame)?; note_snapshot_rendered( @@ -6848,63 +6847,9 @@ async fn run_terminal( .await?; } } - let before_focus_strip = bytes.len(); - bytes = strip_terminal_focus_reports(&bytes); - if before_focus_strip != bytes.len() { - dosh::trace::event( - "client.focus_stripped", - &[ - ("before", before_focus_strip.to_string()), - ("after", bytes.len().to_string()), - ], - ); - } if bytes.is_empty() { continue; } - let before_mouse_strip = bytes.len(); - let (stripped_bytes, stripped_unowned_mouse) = strip_unowned_terminal_reports( - bytes, - predictor.mouse_tracking.enabled(), - ); - bytes = stripped_bytes; - if stripped_unowned_mouse { - dosh::trace::event( - "client.unowned_mouse_stripped", - &[ - ("before", before_mouse_strip.to_string()), - ("after", bytes.len().to_string()), - ("summary", dosh::trace::bytes_summary(&bytes)), - ], - ); - if bytes.is_empty() { - continue; - } - } - if should_strip_stale_terminal_reports( - last_packet_at.elapsed(), - stale_terminal_input_suppress_until, - ) { - let before_mouse_strip = bytes.len(); - bytes = strip_stale_mouse_reports(&bytes); - dosh::trace::event( - "client.stale_strip", - &[ - ("before", before_mouse_strip.to_string()), - ("after", bytes.len().to_string()), - ("silent_ms", last_packet_at.elapsed().as_millis().to_string()), - ( - "grace", - stale_terminal_input_suppress_until - .is_some_and(|deadline| Instant::now() < deadline) - .to_string(), - ), - ], - ); - if bytes.is_empty() { - continue; - } - } if cred.mode != "view-only" && cred.mode != "forward-only" { if let Some(deadline) = startup_input_hold_until { if !predictor.alternate_screen @@ -7047,7 +6992,7 @@ async fn run_terminal( ); refresh_live_addr(&mut addr, &cred)?; arm_stale_terminal_input_suppression( - &mut stale_terminal_input_suppress_until, + &mut stale_terminal_input_suppress_through, ); if !forward_only { render_frame(&frame)?; @@ -7120,7 +7065,7 @@ async fn run_terminal( { refresh_live_addr(&mut addr, &cred)?; arm_stale_terminal_input_suppression( - &mut stale_terminal_input_suppress_until, + &mut stale_terminal_input_suppress_through, ); if !forward_only { render_frame(&frame)?; @@ -7184,7 +7129,7 @@ async fn run_terminal( { refresh_live_addr(&mut addr, &cred)?; arm_stale_terminal_input_suppression( - &mut stale_terminal_input_suppress_until, + &mut stale_terminal_input_suppress_through, ); if !forward_only { render_frame(&frame)?; @@ -7368,7 +7313,7 @@ async fn run_terminal( { refresh_live_addr(&mut addr, &cred)?; arm_stale_terminal_input_suppression( - &mut stale_terminal_input_suppress_until, + &mut stale_terminal_input_suppress_through, ); if !forward_only { render_frame(&frame)?; @@ -7987,7 +7932,7 @@ async fn run_terminal( { refresh_live_addr(&mut addr, &cred)?; arm_stale_terminal_input_suppression( - &mut stale_terminal_input_suppress_until, + &mut stale_terminal_input_suppress_through, ); if !forward_only { render_frame(&frame)?; @@ -8072,7 +8017,7 @@ async fn run_terminal( { refresh_live_addr(&mut addr, &cred)?; arm_stale_terminal_input_suppression( - &mut stale_terminal_input_suppress_until, + &mut stale_terminal_input_suppress_through, ); render_frame(&frame)?; note_snapshot_rendered( @@ -8580,10 +8525,11 @@ fn should_flush_terminal_input_after_contact( fn should_strip_stale_terminal_reports( silent_for: Duration, - suppress_until: Option, + read_at: Instant, + suppress_through: Option, ) -> bool { silent_for >= STALE_TERMINAL_INPUT_AFTER - || suppress_until.is_some_and(|deadline| Instant::now() < deadline) + || suppress_through.is_some_and(|cutoff| read_at <= cutoff) } fn should_reconnect_before_input_for_local_sleep( @@ -8614,44 +8560,217 @@ fn trace_wake_repaint_retry_arm(trigger: &str, elapsed: Duration) { ); } -fn should_strip_unowned_terminal_reports(mouse_tracking: bool) -> bool { - // Mouse mode, not alternate-screen mode, establishes ownership. Inline TUIs - // and multiplexers can legitimately request mouse reports on the main screen. - !mouse_tracking +#[derive(Debug)] +enum TerminalInputEvent { + Raw { bytes: Vec, read_at: Instant }, + Filtered(TerminalInputBatch), } -fn strip_unowned_terminal_reports(bytes: Vec, mouse_tracking: bool) -> (Vec, bool) { - if !should_strip_unowned_terminal_reports(mouse_tracking) { - return (bytes, false); +#[derive(Clone, Copy, Debug)] +struct TerminalReportPolicy { + mouse_tracking: bool, + focus_tracking: bool, + suppress_mouse: bool, +} + +#[derive(Debug)] +struct TerminalInputBatch { + bytes: Vec, + saw_focus_in: bool, + stripped: bool, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum TerminalReportKind { + FocusIn, + FocusOut, + Mouse, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum TerminalReportParse { + Complete { + kind: TerminalReportKind, + len: usize, + }, + Incomplete(Option), + NotReport, +} + +#[derive(Debug, Default)] +struct TerminalInputFilter { + pending: Vec, + pending_since: Option, + pending_policy: Option, +} + +impl TerminalInputFilter { + fn push( + &mut self, + bytes: Vec, + policy: TerminalReportPolicy, + now: Instant, + ) -> TerminalInputBatch { + if self.pending.is_empty() { + self.pending_since = Some(now); + self.pending_policy = Some(policy); + } else if let Some(pending) = &mut self.pending_policy { + pending.mouse_tracking = policy.mouse_tracking; + pending.focus_tracking = policy.focus_tracking; + pending.suppress_mouse |= policy.suppress_mouse; + } + self.pending.extend_from_slice(&bytes); + self.drain(self.pending.len() > MAX_TERMINAL_REPORT_BYTES) } - let before = bytes.len(); - let stripped = strip_complete_terminal_mouse_reports(&bytes); - let changed = stripped.len() != before; - (stripped, changed) -} -fn input_contains_focus_in(bytes: &[u8]) -> bool { - contains_bytes(bytes, b"\x1b[I") || contains_bytes(bytes, b"\x9bI") -} + fn has_pending(&self) -> bool { + !self.pending.is_empty() + } -fn strip_terminal_focus_reports(bytes: &[u8]) -> Vec { - let mut out = Vec::with_capacity(bytes.len()); - let mut offset = 0; - while offset < bytes.len() { - match bytes.get(offset..) { - Some([0x1b, b'[', b'I' | b'O', ..]) => { - offset += 3; - } - Some([0x9b, b'I' | b'O', ..]) => { - offset += 2; - } - _ => { - out.push(bytes[offset]); - offset += 1; + fn flush_due(&mut self, now: Instant) -> Option { + let since = self.pending_since?; + (now.duration_since(since) >= TERMINAL_REPORT_HOLD).then(|| self.drain(true)) + } + + fn drain(&mut self, final_chunk: bool) -> TerminalInputBatch { + let policy = self.pending_policy.unwrap_or(TerminalReportPolicy { + mouse_tracking: false, + focus_tracking: false, + suppress_mouse: false, + }); + let bytes = std::mem::take(&mut self.pending); + let mut output = Vec::with_capacity(bytes.len()); + let mut saw_focus_in = false; + let mut stripped = false; + let mut offset = 0usize; + + while offset < bytes.len() { + match terminal_report_at(&bytes, offset) { + TerminalReportParse::Complete { kind, len } => { + saw_focus_in |= kind == TerminalReportKind::FocusIn; + let drop = match kind { + TerminalReportKind::Mouse => { + policy.suppress_mouse || !policy.mouse_tracking + } + TerminalReportKind::FocusIn | TerminalReportKind::FocusOut => { + !policy.focus_tracking + } + }; + if drop { + stripped = true; + } else { + output.extend_from_slice(&bytes[offset..offset + len]); + } + offset += len; + } + TerminalReportParse::Incomplete(_kind) if !final_chunk => { + self.pending.extend_from_slice(&bytes[offset..]); + self.pending_since.get_or_insert_with(Instant::now); + self.pending_policy = Some(policy); + break; + } + TerminalReportParse::Incomplete(kind) => { + let drop = kind == Some(TerminalReportKind::Mouse) + && (policy.suppress_mouse || !policy.mouse_tracking); + if drop { + stripped = true; + } else { + output.extend_from_slice(&bytes[offset..]); + } + offset = bytes.len(); + } + TerminalReportParse::NotReport => { + output.push(bytes[offset]); + offset += 1; + } } } + + if self.pending.is_empty() { + self.pending_since = None; + self.pending_policy = None; + } + TerminalInputBatch { + bytes: output, + saw_focus_in, + stripped, + } } - out +} + +fn terminal_report_at(bytes: &[u8], offset: usize) -> TerminalReportParse { + let Some(first) = bytes.get(offset).copied() else { + return TerminalReportParse::NotReport; + }; + let (mut cursor, csi_prefix) = match first { + 0x1b => { + let Some(next) = bytes.get(offset + 1).copied() else { + return TerminalReportParse::Incomplete(None); + }; + if next != b'[' { + return TerminalReportParse::NotReport; + } + (offset + 2, 2usize) + } + 0x9b => (offset + 1, 1usize), + _ => return TerminalReportParse::NotReport, + }; + let Some(marker) = bytes.get(cursor).copied() else { + return TerminalReportParse::Incomplete(None); + }; + match marker { + b'I' => TerminalReportParse::Complete { + kind: TerminalReportKind::FocusIn, + len: csi_prefix + 1, + }, + b'O' => TerminalReportParse::Complete { + kind: TerminalReportKind::FocusOut, + len: csi_prefix + 1, + }, + b'M' => { + let len = csi_prefix + 4; + if bytes.len().saturating_sub(offset) < len { + TerminalReportParse::Incomplete(Some(TerminalReportKind::Mouse)) + } else { + TerminalReportParse::Complete { + kind: TerminalReportKind::Mouse, + len, + } + } + } + b'<' => { + cursor += 1; + terminal_parameter_mouse_report(bytes, offset, cursor, true) + } + b'0'..=b'9' => terminal_parameter_mouse_report(bytes, offset, cursor, false), + _ => TerminalReportParse::NotReport, + } +} + +fn terminal_parameter_mouse_report( + bytes: &[u8], + start: usize, + mut cursor: usize, + definite_mouse: bool, +) -> TerminalReportParse { + let mut semicolons = 0usize; + while let Some(byte) = bytes.get(cursor).copied() { + match byte { + b'0'..=b'9' => cursor += 1, + b';' => { + semicolons += 1; + cursor += 1; + } + b'M' | b'm' if semicolons >= 2 => { + return TerminalReportParse::Complete { + kind: TerminalReportKind::Mouse, + len: cursor + 1 - start, + }; + } + _ => return TerminalReportParse::NotReport, + } + } + TerminalReportParse::Incomplete(definite_mouse.then_some(TerminalReportKind::Mouse)) } fn should_repaint_idle_terminal( @@ -8666,14 +8785,12 @@ fn should_repaint_idle_terminal( && now.duration_since(last_attempt_at) >= LOCAL_SLEEP_REPAINT_RETRY_AFTER } -fn arm_stale_terminal_input_suppression(suppress_until: &mut Option) { - *suppress_until = Some(Instant::now() + POST_RECONNECT_STALE_INPUT_GRACE); +fn arm_stale_terminal_input_suppression(suppress_through: &mut Option) { + let cutoff = Instant::now(); + *suppress_through = Some(cutoff); dosh::trace::event( "client.stale_terminal_quarantine_arm", - &[( - "grace_ms", - POST_RECONNECT_STALE_INPUT_GRACE.as_millis().to_string(), - )], + &[("cutoff", "queued-before-reconnect".to_string())], ); } @@ -8755,42 +8872,6 @@ fn strip_stale_mouse_reports(bytes: &[u8]) -> Vec { out } -fn strip_complete_terminal_mouse_reports(bytes: &[u8]) -> Vec { - let mut out = Vec::with_capacity(bytes.len()); - let mut offset = 0; - while offset < bytes.len() { - if let Some(len) = prefixed_mouse_report_len(bytes, offset) { - offset += len; - } else { - out.push(bytes[offset]); - offset += 1; - } - } - out -} - -fn prefixed_mouse_report_len(bytes: &[u8], offset: usize) -> Option { - match bytes.get(offset).copied()? { - 0x1b if bytes.get(offset + 1) == Some(&b'[') => match bytes.get(offset + 2).copied() { - Some(b'M') if offset + 6 <= bytes.len() => Some(6), - Some(b'<') => mouse_report_params_len(bytes, offset + 3, 2).map(|len| len + 3), - Some(byte) if byte.is_ascii_digit() => { - mouse_report_params_len(bytes, offset + 2, 2).map(|len| len + 2) - } - _ => None, - }, - 0x9b => match bytes.get(offset + 1).copied() { - Some(b'M') if offset + 5 <= bytes.len() => Some(5), - Some(b'<') => mouse_report_params_len(bytes, offset + 2, 2).map(|len| len + 2), - Some(byte) if byte.is_ascii_digit() => { - mouse_report_params_len(bytes, offset + 1, 2).map(|len| len + 1) - } - _ => None, - }, - _ => None, - } -} - fn stale_mouse_report_len(bytes: &[u8], offset: usize) -> Option { match bytes.get(offset).copied()? { 0x1b if bytes.get(offset + 1) == Some(&b'[') => match bytes.get(offset + 2).copied() { @@ -9659,15 +9740,17 @@ impl TerminalMouseMode { #[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] struct TerminalModeChanges { alternate_screen: Option, + focus_tracking: Option, mouse_tracking_enable: Option, mouse_tracking_disable_mask: u8, } fn terminal_modes_after_output( mut alternate_screen: bool, + mut focus_tracking: bool, mut mouse_tracking: TerminalMouseMode, bytes: &[u8], -) -> (bool, TerminalMouseMode) { +) -> (bool, bool, TerminalMouseMode) { let mut offset = 0usize; while offset < bytes.len() { let Some((next, changes)) = terminal_private_mode_changes(bytes, offset) else { @@ -9677,6 +9760,9 @@ fn terminal_modes_after_output( if let Some(enable) = changes.alternate_screen { alternate_screen = enable; } + if let Some(enable) = changes.focus_tracking { + focus_tracking = enable; + } if let Some(mode) = changes.mouse_tracking_enable { mouse_tracking = mode; } else if terminal_mouse_mode_mask(mouse_tracking) & changes.mouse_tracking_disable_mask @@ -9686,7 +9772,7 @@ fn terminal_modes_after_output( } offset = next.max(offset + 1); } - (alternate_screen, mouse_tracking) + (alternate_screen, focus_tracking, mouse_tracking) } #[cfg(test)] @@ -9725,6 +9811,8 @@ fn terminal_private_mode_changes( TerminalModeChanges { alternate_screen: terminal_private_params_include_alt_screen(params) .then_some(enable), + focus_tracking: terminal_private_params_include_focus_tracking(params) + .then_some(enable), mouse_tracking_enable, mouse_tracking_disable_mask, }, @@ -9743,6 +9831,12 @@ fn terminal_private_params_include_alt_screen(params: &[u8]) -> bool { .any(|param| matches!(param, b"47" | b"1047" | b"1049")) } +fn terminal_private_params_include_focus_tracking(params: &[u8]) -> bool { + params + .split(|byte| matches!(byte, b';' | b':')) + .any(|param| param == b"1004") +} + fn terminal_private_mouse_tracking_changes( params: &[u8], enable: bool, @@ -9808,6 +9902,8 @@ struct Predictor { /// When false, SGR mouse bytes from the local terminal are stale UI noise and /// must not be forwarded into the shell prompt. mouse_tracking: TerminalMouseMode, + /// True while the remote program has requested terminal focus reports. + focus_tracking: bool, /// Tail of recent terminal output, retained so alternate-screen transitions /// split across UDP frames are still detected. output_parse_tail: Vec, @@ -9859,6 +9955,7 @@ impl Predictor { output_backpressured: false, alternate_screen: false, mouse_tracking: TerminalMouseMode::None, + focus_tracking: false, output_parse_tail: Vec::new(), cells: Vec::new(), cursor: 0, @@ -9883,6 +9980,7 @@ impl Predictor { self.confirmed_epoch = self.epoch - 1; self.alternate_screen = false; self.mouse_tracking = TerminalMouseMode::None; + self.focus_tracking = false; self.output_parse_tail.clear(); self.oldest_pending_at = None; } @@ -10046,20 +10144,28 @@ impl Predictor { let before_alternate_screen = self.alternate_screen; let before_mouse_tracking = self.mouse_tracking; - let (alternate_screen, mouse_tracking) = - terminal_modes_after_output(self.alternate_screen, self.mouse_tracking, &parse); + let before_focus_tracking = self.focus_tracking; + let (alternate_screen, focus_tracking, mouse_tracking) = terminal_modes_after_output( + self.alternate_screen, + self.focus_tracking, + self.mouse_tracking, + &parse, + ); self.alternate_screen = alternate_screen; + self.focus_tracking = focus_tracking; self.mouse_tracking = mouse_tracking; if self.alternate_screen != before_alternate_screen { let _ = self.discard_all(); } if self.alternate_screen != before_alternate_screen + || self.focus_tracking != before_focus_tracking || self.mouse_tracking != before_mouse_tracking { dosh::trace::event( "client.terminal_modes", &[ ("alt", self.alternate_screen.to_string()), + ("focus", self.focus_tracking.to_string()), ("mouse", self.mouse_tracking.enabled().to_string()), ("bytes", dosh::trace::bytes_summary(bytes)), ], @@ -11191,15 +11297,15 @@ mod tests { CachedCredential, DisconnectStatus, DynamicForward, FRAME_GAP_RESYNC_AFTER_MS, FrameBuffer, LOCAL_SLEEP_REPAINT_AFTER, LOCAL_SLEEP_REPAINT_RETRY_AFTER, LOCAL_SLEEP_REPAINT_RETRY_WINDOW, LocalForward, MAX_PENDING_USER_INPUT_BYTES, - NativeIdentityContext, POST_RECONNECT_STALE_INPUT_GRACE, POST_SUBMIT_ALL_INPUT_HOLD, - PendingStreamControl, PendingStreamOpen, PendingWindowAdjust, PredictMode, Predictor, - RESTART_STATUS_SCRIPT, RemoteForward, STALE_TERMINAL_INPUT_AFTER, STARTUP_INPUT_HOLD, - STDIN_QUEUE_CAPACITY, STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS, STREAM_INITIAL_WINDOW, - SshConfig, SshPathTokenContext, StartupGateMode, StatusAction, TERMINAL_CLEANUP, - TERMINAL_SNAPSHOT_RESET, UpdateOptions, UpdateRole, auth_allows, cache_key, - cache_server_prefix, cleanup_stream_state, clear_cached_credentials, - ensure_tui_safe_status_overlay, expand_ssh_path_tokens, first_resolved_addr, - imported_host_block, input_contains_focus_in, input_prefix_before_escape, + NativeIdentityContext, POST_SUBMIT_ALL_INPUT_HOLD, PendingStreamControl, PendingStreamOpen, + PendingWindowAdjust, PredictMode, Predictor, RESTART_STATUS_SCRIPT, RemoteForward, + STALE_TERMINAL_INPUT_AFTER, STARTUP_INPUT_HOLD, STDIN_QUEUE_CAPACITY, + STREAM_CONTROL_RETRANSMIT_MAX_ATTEMPTS, STREAM_INITIAL_WINDOW, SshConfig, + SshPathTokenContext, StartupGateMode, StatusAction, TERMINAL_CLEANUP, + TERMINAL_SNAPSHOT_RESET, TerminalInputFilter, TerminalReportPolicy, UpdateOptions, + UpdateRole, auth_allows, cache_key, cache_server_prefix, cleanup_stream_state, + clear_cached_credentials, ensure_tui_safe_status_overlay, expand_ssh_path_tokens, + first_resolved_addr, imported_host_block, input_prefix_before_escape, is_local_status_target, is_resume_response_for_client, load_first_native_identity_with_prompt, local_symlink_target_is_dir, local_username_from_env, native_proxy_udp_warning, newest_client_trace_path_from, @@ -11217,20 +11323,19 @@ mod tests { send_stream_eof, server_version_mismatch, should_flush_terminal_input_after_contact, should_health_log_client_start, should_hold_during_startup_gate, should_hold_post_submit_input, should_reconnect_before_input_for_local_sleep, - should_repaint_idle_terminal, should_strip_unowned_terminal_reports, socket_bind_display, + should_repaint_idle_terminal, should_strip_stale_terminal_reports, socket_bind_display, split_after_command_submit, split_forward_spec, split_trace_tokens, ssh_command_target, ssh_config_uses_proxy, ssh_config_word_for_os, ssh_destination_host, ssh_username, ssh_with_user, startup_command, status_ssh_target, strip_stale_mouse_reports, - strip_terminal_focus_reports, strip_unowned_terminal_reports, summarize_trace_file, - summarize_trace_file_with_mode, terminal_input_channel, terminal_private_mode_transition, - terminal_resize_poll_interval_for, toml_bare_key_or_quoted, top_trace_events, - trace_report_warnings, unix_update_script, update_installer_url, update_version_status, - upsert_managed_block, valid_forward_host, vscode_command_candidates, - vscode_fallback_command, vscode_safe_alias, wake_repaint_retry_deadline, - windows_command_word, windows_deferred_update_script, windows_mode_from_readonly, - windows_powershell_command_candidates, windows_readonly_from_mode, windows_update_script, - windows_url_contents_script, windows_url_reachable_script, windows_vt_input_mode, - windows_vt_output_mode, + summarize_trace_file, summarize_trace_file_with_mode, terminal_input_channel, + terminal_private_mode_transition, terminal_resize_poll_interval_for, + toml_bare_key_or_quoted, top_trace_events, trace_report_warnings, unix_update_script, + update_installer_url, update_version_status, upsert_managed_block, valid_forward_host, + vscode_command_candidates, vscode_fallback_command, vscode_safe_alias, + wake_repaint_retry_deadline, windows_command_word, windows_deferred_update_script, + windows_mode_from_readonly, windows_powershell_command_candidates, + windows_readonly_from_mode, windows_update_script, windows_url_contents_script, + windows_url_reachable_script, windows_vt_input_mode, windows_vt_output_mode, }; use dosh::config::{ClientConfig, CommandExtension, HostConfig}; use dosh::native::EnvVar; @@ -11241,6 +11346,49 @@ mod tests { use std::fs; use std::time::{Duration, Instant}; + fn terminal_report_policy(mouse_tracking: bool, focus_tracking: bool) -> TerminalReportPolicy { + TerminalReportPolicy { + mouse_tracking, + focus_tracking, + suppress_mouse: false, + } + } + + fn filter_terminal_parts( + parts: &[&[u8]], + policy: TerminalReportPolicy, + ) -> (Vec, bool, bool) { + let mut filter = TerminalInputFilter::default(); + let mut output = Vec::new(); + let mut saw_focus_in = false; + let mut stripped = false; + for part in parts { + let batch = filter.push(part.to_vec(), policy, Instant::now()); + output.extend_from_slice(&batch.bytes); + saw_focus_in |= batch.saw_focus_in; + stripped |= batch.stripped; + } + let batch = filter.drain(true); + output.extend_from_slice(&batch.bytes); + saw_focus_in |= batch.saw_focus_in; + stripped |= batch.stripped; + (output, saw_focus_in, stripped) + } + + fn assert_terminal_filter_at_every_split( + input: &[u8], + expected: &[u8], + policy: TerminalReportPolicy, + ) { + for split in 1..input.len() { + let (actual, _, _) = filter_terminal_parts(&[&input[..split], &input[split..]], policy); + assert_eq!( + actual, expected, + "terminal input split at {split}: {input:?}" + ); + } + } + fn test_args(server: &str, command: &[&str]) -> super::Args { super::Args { server: Some(server.to_string()), @@ -12310,22 +12458,40 @@ mod tests { #[test] fn main_screen_tui_mouse_input_survives_until_mode_is_disabled() { let mut predictor = Predictor::new(true); - let report = b"\x1b[<0;12;8M".to_vec(); + let report = b"\x1b[<0;12;8M"; predictor.observe_output(b"\x1b[?1000;1006h"); assert!(!predictor.alternate_screen); - let (preserved, changed) = - strip_unowned_terminal_reports(report.clone(), predictor.mouse_tracking.enabled()); + let (preserved, _, changed) = filter_terminal_parts( + &[report], + terminal_report_policy(predictor.mouse_tracking.enabled(), predictor.focus_tracking), + ); assert!(!changed); assert_eq!(preserved, report); predictor.observe_output(b"\x1b[?1006;1000l"); - let (stripped, changed) = - strip_unowned_terminal_reports(report, predictor.mouse_tracking.enabled()); + let (stripped, _, changed) = filter_terminal_parts( + &[report], + terminal_report_policy(predictor.mouse_tracking.enabled(), predictor.focus_tracking), + ); assert!(changed); assert!(stripped.is_empty()); } + #[test] + fn focus_tracking_is_independent_and_survives_split_mode_output() { + let mut predictor = Predictor::new(true); + + predictor.observe_output(b"\x1b[?10"); + predictor.observe_output(b"04h"); + assert!(predictor.focus_tracking); + assert!(!predictor.mouse_tracking.enabled()); + + predictor.observe_output(b"\x1b[?1000h\x1b[?1004l"); + assert!(!predictor.focus_tracking); + assert!(predictor.mouse_tracking.enabled()); + } + #[test] fn mouse_encoding_alone_does_not_claim_input_and_cannot_disable_tracking() { let mut predictor = Predictor::new(true); @@ -12390,9 +12556,17 @@ mod tests { } #[test] - fn mouse_tracking_alone_establishes_terminal_report_ownership() { - assert!(should_strip_unowned_terminal_reports(false)); - assert!(!should_strip_unowned_terminal_reports(true)); + fn mouse_tracking_establishes_terminal_report_ownership() { + let report = b"\x1b[<35;10;1M"; + let (stripped, _, changed) = + filter_terminal_parts(&[report], terminal_report_policy(false, false)); + assert!(changed); + assert!(stripped.is_empty()); + + let (preserved, _, changed) = + filter_terminal_parts(&[report], terminal_report_policy(true, false)); + assert!(!changed); + assert_eq!(preserved, report); } #[test] @@ -14304,19 +14478,21 @@ mod tests { #[test] fn unowned_terminal_mouse_reports_strip_unless_application_owns_mouse() { - let input = b"\x1b[<35;10;1Mcmd\r".to_vec(); + let input = b"\x1b[<35;10;1Mcmd\r"; - let (stripped, changed) = strip_unowned_terminal_reports(input.clone(), false); + let (stripped, _, changed) = + filter_terminal_parts(&[input], terminal_report_policy(false, false)); assert!(changed); assert_eq!(stripped, b"cmd\r"); - let (preserved, changed) = strip_unowned_terminal_reports(b"\x1b[<35;10;1M".to_vec(), true); + let (preserved, _, changed) = + filter_terminal_parts(&[b"\x1b[<35;10;1M"], terminal_report_policy(true, false)); assert!(!changed); assert_eq!(preserved, b"\x1b[<35;10;1M"); } #[test] - fn live_unowned_filter_preserves_ambiguous_user_input() { + fn stateful_report_filter_preserves_non_report_input_at_every_split() { for input in [ b"\x1b".as_slice(), b"\x1b[".as_slice(), @@ -14325,9 +14501,42 @@ mod tests { b"35;152;1M".as_slice(), b"printf 35;152;1M\r".as_slice(), ] { - let (preserved, changed) = strip_unowned_terminal_reports(input.to_vec(), false); + let (preserved, _, changed) = + filter_terminal_parts(&[input], terminal_report_policy(false, false)); assert!(!changed, "live filter changed {input:?}"); assert_eq!(preserved, input); + assert_terminal_filter_at_every_split( + input, + input, + terminal_report_policy(false, false), + ); + } + } + + #[test] + fn stateful_report_filter_frames_every_mouse_protocol_at_every_split() { + let reports = [ + b"\x1b[<35;152;1M".as_slice(), + b"\x1b[35;152;1M".as_slice(), + b"\x1b[M !!".as_slice(), + b"\x9b<35;152;1m".as_slice(), + b"\x9b35;152;1M".as_slice(), + b"\x9bM !!".as_slice(), + ]; + for report in reports { + let mut framed = b"before".to_vec(); + framed.extend_from_slice(report); + framed.extend_from_slice(b"after"); + assert_terminal_filter_at_every_split( + &framed, + b"beforeafter", + terminal_report_policy(false, false), + ); + assert_terminal_filter_at_every_split( + &framed, + &framed, + terminal_report_policy(true, false), + ); } } @@ -14441,22 +14650,31 @@ mod tests { } #[test] - fn focus_in_reports_trigger_local_repaint_detection() { - assert!(input_contains_focus_in(b"\x1b[I")); - assert!(input_contains_focus_in(b"before\x9bIafter")); - assert!(!input_contains_focus_in(b"\x1b[O")); - assert!(!input_contains_focus_in(b"\x1b[A")); - } + fn focus_reports_are_mode_aware_and_framed_at_every_split() { + for report in [b"\x1b[I".as_slice(), b"\x1b[O", b"\x9bI", b"\x9bO"] { + let saw_focus_in = report.ends_with(b"I"); + let (stripped, saw, changed) = + filter_terminal_parts(&[report], terminal_report_policy(false, false)); + assert!(changed); + assert!(stripped.is_empty()); + assert_eq!(saw, saw_focus_in); + assert_terminal_filter_at_every_split( + report, + b"", + terminal_report_policy(false, false), + ); - #[test] - fn focus_reports_are_never_forwarded_as_user_input() { - assert_eq!(strip_terminal_focus_reports(b"\x1b[Ihello\x1b[O"), b"hello"); - assert_eq!(strip_terminal_focus_reports(b"\x9bIhello\x9bO"), b"hello"); - assert_eq!(strip_terminal_focus_reports(b"\x1b[A"), b"\x1b[A"); - assert_eq!( - strip_terminal_focus_reports(b"before\x1b[Iafter"), - b"beforeafter" - ); + let (preserved, saw, changed) = + filter_terminal_parts(&[report], terminal_report_policy(false, true)); + assert!(!changed); + assert_eq!(preserved, report); + assert_eq!(saw, saw_focus_in); + assert_terminal_filter_at_every_split( + report, + report, + terminal_report_policy(false, true), + ); + } } #[test] @@ -14557,12 +14775,34 @@ mod tests { } #[test] - fn reconnect_mouse_quarantine_is_long_enough_for_sleep_wake_noise() { - assert!(POST_RECONNECT_STALE_INPUT_GRACE >= LOCAL_SLEEP_REPAINT_AFTER); + fn reconnect_mouse_quarantine_strips_queued_noise_without_a_time_grace() { assert_eq!( strip_stale_mouse_reports(b"35;152;1M\x1b[A35;149;1M"), b"\x1b[A" ); + + let cutoff = Instant::now(); + assert!(should_strip_stale_terminal_reports( + Duration::ZERO, + cutoff - Duration::from_millis(1), + Some(cutoff) + )); + assert!(!should_strip_stale_terminal_reports( + Duration::ZERO, + cutoff + Duration::from_millis(1), + Some(cutoff) + )); + + let (filtered, _, stripped) = filter_terminal_parts( + &[b"\x1b[<35;10;1Mkey"], + TerminalReportPolicy { + mouse_tracking: true, + focus_tracking: false, + suppress_mouse: true, + }, + ); + assert!(stripped); + assert_eq!(filtered, b"key"); } #[test] diff --git a/tests/client_terminal_runtime.rs b/tests/client_terminal_runtime.rs index 42e1a99..b7ef85b 100644 --- a/tests/client_terminal_runtime.rs +++ b/tests/client_terminal_runtime.rs @@ -16,6 +16,8 @@ const SESSION_KEY: [u8; 32] = [0x91; 32]; const INPUT_BYTES: &[u8] = concat!( "\x1b[A", "\x1b[B", + "\x1b[I", + "\x1b[O", "\x1b[200~", "paste-λ-界", "\x1b[201~", @@ -149,6 +151,10 @@ fn native_client_terminal_round_trip_is_platform_complete() { contains(&output, b"\x1b[?1006h"), "SGR mouse mode was not preserved: {output:?}" ); + assert!( + contains(&output, b"\x1b[?1004h"), + "focus-report mode was not preserved: {output:?}" + ); assert!( contains(&output, b"\x1b[?25h"), "terminal cleanup did not restore the cursor: {output:?}" @@ -565,6 +571,7 @@ fn run_fake_terminal_server(socket: UdpSocket, observations: mpsc::Sender