Compare commits
7
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f8693f08b5 | ||
|
|
26532fc0e1 | ||
|
|
c5f699a6ef | ||
|
|
60403ba4c3 | ||
|
|
833ac1082f | ||
|
|
97cf165527 | ||
|
|
8f2d57d95e |
Generated
+1
-1
@@ -436,7 +436,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "dosh"
|
name = "dosh"
|
||||||
version = "1.0.0-rc46"
|
version = "1.0.0-rc49"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"base64",
|
"base64",
|
||||||
|
|||||||
+1
-1
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "dosh"
|
name = "dosh"
|
||||||
version = "1.0.0-rc46"
|
version = "1.0.0-rc49"
|
||||||
edition = "2024"
|
edition = "2024"
|
||||||
license = "MIT"
|
license = "MIT"
|
||||||
|
|
||||||
|
|||||||
+80
-104
@@ -77,7 +77,6 @@ const POST_SUBMIT_ALL_INPUT_HOLD: Duration = Duration::from_millis(120);
|
|||||||
const STALE_TERMINAL_INPUT_AFTER: Duration = Duration::from_secs(2);
|
const STALE_TERMINAL_INPUT_AFTER: Duration = Duration::from_secs(2);
|
||||||
const POST_RECONNECT_STALE_INPUT_GRACE: Duration = Duration::from_secs(5);
|
const POST_RECONNECT_STALE_INPUT_GRACE: Duration = Duration::from_secs(5);
|
||||||
const FOCUS_REPAINT_COOLDOWN: Duration = Duration::from_secs(1);
|
const FOCUS_REPAINT_COOLDOWN: Duration = Duration::from_secs(1);
|
||||||
const ALT_SCREEN_IDLE_REPAINT_AFTER: Duration = Duration::from_secs(15);
|
|
||||||
const LOCAL_SLEEP_REPAINT_AFTER: Duration = Duration::from_secs(5);
|
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_AFTER: Duration = Duration::from_secs(1);
|
||||||
const LOCAL_SLEEP_REPAINT_RETRY_WINDOW: Duration = Duration::from_secs(10);
|
const LOCAL_SLEEP_REPAINT_RETRY_WINDOW: Duration = Duration::from_secs(10);
|
||||||
@@ -6570,6 +6569,8 @@ async fn run_terminal(
|
|||||||
predict_mode,
|
predict_mode,
|
||||||
);
|
);
|
||||||
let mut frame_renderer = TerminalFrameRenderer::new()?;
|
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
|
// Non-destructive disconnect status line. Off in forward-only mode (no TTY to
|
||||||
// draw on) and respecting the client config (default on, env override).
|
// draw on) and respecting the client config (default on, env override).
|
||||||
let mut disconnect_status = DisconnectStatus::new(resolve_disconnect_status() && !forward_only);
|
let mut disconnect_status = DisconnectStatus::new(resolve_disconnect_status() && !forward_only);
|
||||||
@@ -6612,9 +6613,8 @@ async fn run_terminal(
|
|||||||
let mut startup_input_hold_until: Option<Instant> = None;
|
let mut startup_input_hold_until: Option<Instant> = None;
|
||||||
let mut startup_gate_mode = StartupGateMode::HoldControl;
|
let mut startup_gate_mode = StartupGateMode::HoldControl;
|
||||||
let mut stale_terminal_input_suppress_until: Option<Instant> = None;
|
let mut stale_terminal_input_suppress_until: Option<Instant> = None;
|
||||||
let mut last_terminal_frame_at = Instant::now();
|
|
||||||
let mut last_focus_repaint_at = Instant::now() - FOCUS_REPAINT_COOLDOWN;
|
let mut last_focus_repaint_at = Instant::now() - FOCUS_REPAINT_COOLDOWN;
|
||||||
let mut last_idle_repaint_attempt_at = Instant::now() - ALT_SCREEN_IDLE_REPAINT_AFTER;
|
let mut last_idle_repaint_attempt_at = Instant::now() - LOCAL_SLEEP_REPAINT_RETRY_AFTER;
|
||||||
let mut last_status_tick_at = Instant::now();
|
let mut last_status_tick_at = Instant::now();
|
||||||
let mut wake_repaint_retry_until: Option<Instant> = None;
|
let mut wake_repaint_retry_until: Option<Instant> = None;
|
||||||
if let Some(frame) = first_frame {
|
if let Some(frame) = first_frame {
|
||||||
@@ -6622,7 +6622,6 @@ async fn run_terminal(
|
|||||||
render_frame(&frame)?;
|
render_frame(&frame)?;
|
||||||
note_snapshot_rendered(&frame, &mut disconnect_status, &mut status_restore_pending);
|
note_snapshot_rendered(&frame, &mut disconnect_status, &mut status_restore_pending);
|
||||||
predictor.observe_output(&frame.bytes);
|
predictor.observe_output(&frame.bytes);
|
||||||
last_terminal_frame_at = Instant::now();
|
|
||||||
}
|
}
|
||||||
if frame.closed {
|
if frame.closed {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
@@ -6676,26 +6675,37 @@ async fn run_terminal(
|
|||||||
let mut detach_requested = false;
|
let mut detach_requested = false;
|
||||||
loop {
|
loop {
|
||||||
accepted_output_seq = accepted_output_seq.max(cred.last_rendered_seq);
|
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 {
|
if detach_requested {
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
rendered = frame_renderer.recv(), if frame_renderer.has_pending() => {
|
rendered = frame_renderer.recv(), if frame_renderer.has_pending() => {
|
||||||
let frame = rendered?;
|
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);
|
cred.last_rendered_seq = cred.last_rendered_seq.max(frame.output_seq);
|
||||||
note_snapshot_rendered(
|
note_snapshot_rendered(
|
||||||
&frame,
|
&frame,
|
||||||
&mut disconnect_status,
|
&mut disconnect_status,
|
||||||
&mut status_restore_pending,
|
&mut status_restore_pending,
|
||||||
);
|
);
|
||||||
last_terminal_frame_at = Instant::now();
|
|
||||||
wake_repaint_retry_until = None;
|
wake_repaint_retry_until = None;
|
||||||
send_ack(&socket, addr, &cred, &mut send_seq).await?;
|
send_ack(&socket, addr, &cred, &mut send_seq).await?;
|
||||||
if frame.closed {
|
if frame.closed {
|
||||||
return Ok(());
|
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() => {
|
stdin_msg = stdin_rx.recv() => {
|
||||||
match stdin_msg {
|
match stdin_msg {
|
||||||
@@ -6765,7 +6775,6 @@ async fn run_terminal(
|
|||||||
&mut status_restore_pending,
|
&mut status_restore_pending,
|
||||||
);
|
);
|
||||||
predictor.observe_output(&frame.bytes);
|
predictor.observe_output(&frame.bytes);
|
||||||
last_terminal_frame_at = Instant::now();
|
|
||||||
last_packet_at = Instant::now();
|
last_packet_at = Instant::now();
|
||||||
last_focus_repaint_at = Instant::now();
|
last_focus_repaint_at = Instant::now();
|
||||||
wake_repaint_retry_until = None;
|
wake_repaint_retry_until = None;
|
||||||
@@ -6825,7 +6834,6 @@ async fn run_terminal(
|
|||||||
&mut status_restore_pending,
|
&mut status_restore_pending,
|
||||||
);
|
);
|
||||||
predictor.observe_output(&frame.bytes);
|
predictor.observe_output(&frame.bytes);
|
||||||
last_terminal_frame_at = Instant::now();
|
|
||||||
last_packet_at = Instant::now();
|
last_packet_at = Instant::now();
|
||||||
wake_repaint_retry_until = None;
|
wake_repaint_retry_until = None;
|
||||||
flush_pending_user_input(
|
flush_pending_user_input(
|
||||||
@@ -7049,7 +7057,6 @@ async fn run_terminal(
|
|||||||
&mut status_restore_pending,
|
&mut status_restore_pending,
|
||||||
);
|
);
|
||||||
predictor.observe_output(&frame.bytes);
|
predictor.observe_output(&frame.bytes);
|
||||||
last_terminal_frame_at = Instant::now();
|
|
||||||
}
|
}
|
||||||
last_packet_at = Instant::now();
|
last_packet_at = Instant::now();
|
||||||
flush_pending_user_input(
|
flush_pending_user_input(
|
||||||
@@ -7098,7 +7105,8 @@ async fn run_terminal(
|
|||||||
maybe_send_resize(&socket, addr, &cred, &mut send_seq, &mut last_size).await?;
|
maybe_send_resize(&socket, addr, &cred, &mut send_seq, &mut last_size).await?;
|
||||||
}
|
}
|
||||||
_ = frame_gap_tick.tick() => {
|
_ = frame_gap_tick.tick() => {
|
||||||
if !frame_renderer.has_pending()
|
if !render_resync_needed
|
||||||
|
&& !frame_renderer.has_pending()
|
||||||
&& frame_buffer.resync_due()
|
&& frame_buffer.resync_due()
|
||||||
&& let Some(frame) = reconnect(
|
&& let Some(frame) = reconnect(
|
||||||
&socket,
|
&socket,
|
||||||
@@ -7122,7 +7130,6 @@ async fn run_terminal(
|
|||||||
&mut status_restore_pending,
|
&mut status_restore_pending,
|
||||||
);
|
);
|
||||||
predictor.observe_output(&frame.bytes);
|
predictor.observe_output(&frame.bytes);
|
||||||
last_terminal_frame_at = Instant::now();
|
|
||||||
wake_repaint_retry_until = None;
|
wake_repaint_retry_until = None;
|
||||||
}
|
}
|
||||||
last_packet_at = Instant::now();
|
last_packet_at = Instant::now();
|
||||||
@@ -7187,7 +7194,6 @@ async fn run_terminal(
|
|||||||
&mut status_restore_pending,
|
&mut status_restore_pending,
|
||||||
);
|
);
|
||||||
predictor.observe_output(&frame.bytes);
|
predictor.observe_output(&frame.bytes);
|
||||||
last_terminal_frame_at = Instant::now();
|
|
||||||
wake_repaint_retry_until = None;
|
wake_repaint_retry_until = None;
|
||||||
}
|
}
|
||||||
last_packet_at = Instant::now();
|
last_packet_at = Instant::now();
|
||||||
@@ -7217,10 +7223,31 @@ async fn run_terminal(
|
|||||||
predictor.clear_pending()?;
|
predictor.clear_pending()?;
|
||||||
if !forward_only {
|
if !forward_only {
|
||||||
predictor.observe_output(&frame.bytes);
|
predictor.observe_output(&frame.bytes);
|
||||||
last_terminal_frame_at = Instant::now();
|
|
||||||
wake_repaint_retry_until = None;
|
wake_repaint_retry_until = None;
|
||||||
frame_renderer.enqueue(frame)?;
|
if render_resync_needed {
|
||||||
predictor.set_output_backpressured(true);
|
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(
|
flush_startup_input_if_ready(
|
||||||
&socket,
|
&socket,
|
||||||
addr,
|
addr,
|
||||||
@@ -7351,7 +7378,6 @@ async fn run_terminal(
|
|||||||
&mut status_restore_pending,
|
&mut status_restore_pending,
|
||||||
);
|
);
|
||||||
predictor.observe_output(&frame.bytes);
|
predictor.observe_output(&frame.bytes);
|
||||||
last_terminal_frame_at = Instant::now();
|
|
||||||
wake_repaint_retry_until = None;
|
wake_repaint_retry_until = None;
|
||||||
}
|
}
|
||||||
last_packet_at = Instant::now();
|
last_packet_at = Instant::now();
|
||||||
@@ -7945,7 +7971,8 @@ async fn run_terminal(
|
|||||||
let mut repainted_this_tick = false;
|
let mut repainted_this_tick = false;
|
||||||
let stale = last_packet_at.elapsed();
|
let stale = last_packet_at.elapsed();
|
||||||
if !frame_renderer.has_pending()
|
if !frame_renderer.has_pending()
|
||||||
&& (status_restore_pending
|
&& (render_resync_needed
|
||||||
|
|| status_restore_pending
|
||||||
|| stale >= Duration::from_secs(reconnect_timeout_secs.max(1)))
|
|| stale >= Duration::from_secs(reconnect_timeout_secs.max(1)))
|
||||||
{
|
{
|
||||||
if let Some(frame) = reconnect(
|
if let Some(frame) = reconnect(
|
||||||
@@ -7970,10 +7997,10 @@ async fn run_terminal(
|
|||||||
&mut status_restore_pending,
|
&mut status_restore_pending,
|
||||||
);
|
);
|
||||||
predictor.observe_output(&frame.bytes);
|
predictor.observe_output(&frame.bytes);
|
||||||
last_terminal_frame_at = Instant::now();
|
|
||||||
last_idle_repaint_attempt_at = Instant::now();
|
last_idle_repaint_attempt_at = Instant::now();
|
||||||
wake_repaint_retry_until = None;
|
wake_repaint_retry_until = None;
|
||||||
repainted_this_tick = true;
|
repainted_this_tick = true;
|
||||||
|
render_resync_needed = false;
|
||||||
}
|
}
|
||||||
last_packet_at = Instant::now();
|
last_packet_at = Instant::now();
|
||||||
flush_pending_user_input(
|
flush_pending_user_input(
|
||||||
@@ -8004,7 +8031,7 @@ async fn run_terminal(
|
|||||||
// a latency spike (or recovery) flips speculation on/off promptly
|
// a latency spike (or recovery) flips speculation on/off promptly
|
||||||
// without waiting for the next keystroke to drive `redraw`.
|
// without waiting for the next keystroke to drive `redraw`.
|
||||||
if !forward_only {
|
if !forward_only {
|
||||||
if !frame_renderer.has_pending() {
|
if !frame_renderer.has_pending() && !render_resync_needed {
|
||||||
predictor.refresh_policy()?;
|
predictor.refresh_policy()?;
|
||||||
}
|
}
|
||||||
flush_startup_input_if_ready(
|
flush_startup_input_if_ready(
|
||||||
@@ -8018,9 +8045,7 @@ async fn run_terminal(
|
|||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
let now = Instant::now();
|
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,
|
last_idle_repaint_attempt_at,
|
||||||
status_tick_gap,
|
status_tick_gap,
|
||||||
wake_repaint_retry_until,
|
wake_repaint_retry_until,
|
||||||
@@ -8056,7 +8081,6 @@ async fn run_terminal(
|
|||||||
&mut status_restore_pending,
|
&mut status_restore_pending,
|
||||||
);
|
);
|
||||||
predictor.observe_output(&frame.bytes);
|
predictor.observe_output(&frame.bytes);
|
||||||
last_terminal_frame_at = Instant::now();
|
|
||||||
last_packet_at = Instant::now();
|
last_packet_at = Instant::now();
|
||||||
wake_repaint_retry_until = None;
|
wake_repaint_retry_until = None;
|
||||||
flush_pending_user_input(
|
flush_pending_user_input(
|
||||||
@@ -8076,7 +8100,10 @@ async fn run_terminal(
|
|||||||
// on how long the link has been silent (recomputed after any
|
// on how long the link has been silent (recomputed after any
|
||||||
// reconnect attempt above may have reset `last_packet_at`).
|
// reconnect attempt above may have reset `last_packet_at`).
|
||||||
if !forward_only {
|
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()
|
disconnect_status.on_suppressed()
|
||||||
} else {
|
} else {
|
||||||
disconnect_status.on_tick(last_packet_at.elapsed())
|
disconnect_status.on_tick(last_packet_at.elapsed())
|
||||||
@@ -8628,8 +8655,6 @@ fn strip_terminal_focus_reports(bytes: &[u8]) -> Vec<u8> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn should_repaint_idle_terminal(
|
fn should_repaint_idle_terminal(
|
||||||
alternate_screen: bool,
|
|
||||||
last_terminal_frame_at: Instant,
|
|
||||||
last_attempt_at: Instant,
|
last_attempt_at: Instant,
|
||||||
status_tick_gap: Duration,
|
status_tick_gap: Duration,
|
||||||
wake_repaint_retry_until: Option<Instant>,
|
wake_repaint_retry_until: Option<Instant>,
|
||||||
@@ -8637,12 +8662,8 @@ fn should_repaint_idle_terminal(
|
|||||||
) -> bool {
|
) -> bool {
|
||||||
let sleep_wake_gap = status_tick_gap >= LOCAL_SLEEP_REPAINT_AFTER;
|
let sleep_wake_gap = status_tick_gap >= LOCAL_SLEEP_REPAINT_AFTER;
|
||||||
let wake_retry_active = wake_repaint_retry_until.is_some_and(|deadline| now < deadline);
|
let wake_retry_active = wake_repaint_retry_until.is_some_and(|deadline| now < deadline);
|
||||||
let stale_alternate_screen = alternate_screen
|
(sleep_wake_gap || wake_retry_active)
|
||||||
&& now.duration_since(last_terminal_frame_at) >= ALT_SCREEN_IDLE_REPAINT_AFTER;
|
&& now.duration_since(last_attempt_at) >= LOCAL_SLEEP_REPAINT_RETRY_AFTER
|
||||||
if sleep_wake_gap || wake_retry_active {
|
|
||||||
return now.duration_since(last_attempt_at) >= LOCAL_SLEEP_REPAINT_RETRY_AFTER;
|
|
||||||
}
|
|
||||||
stale_alternate_screen && now.duration_since(last_attempt_at) >= ALT_SCREEN_IDLE_REPAINT_AFTER
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn arm_stale_terminal_input_suppression(suppress_until: &mut Option<Instant>) {
|
fn arm_stale_terminal_input_suppression(suppress_until: &mut Option<Instant>) {
|
||||||
@@ -10576,19 +10597,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 {
|
let Some(jobs) = self.jobs.as_ref() else {
|
||||||
return Err(anyhow!("terminal renderer is closed"));
|
return Err(anyhow!("terminal renderer is closed"));
|
||||||
};
|
};
|
||||||
jobs.try_send(frame).map_err(|err| match err {
|
match jobs.try_send(frame) {
|
||||||
mpsc::error::TrySendError::Full(_) => anyhow!(
|
Ok(()) => {
|
||||||
"terminal renderer exceeded its bounded {}-frame queue",
|
self.pending += 1;
|
||||||
TERMINAL_RENDER_QUEUE_CAPACITY
|
Ok(None)
|
||||||
),
|
}
|
||||||
mpsc::error::TrySendError::Closed(_) => anyhow!("terminal renderer stopped"),
|
Err(mpsc::error::TrySendError::Full(frame)) => Ok(Some(frame)),
|
||||||
})?;
|
Err(mpsc::error::TrySendError::Closed(_)) => Err(anyhow!("terminal renderer stopped")),
|
||||||
self.pending += 1;
|
}
|
||||||
Ok(())
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn has_pending(&self) -> bool {
|
fn has_pending(&self) -> bool {
|
||||||
@@ -10652,9 +10676,12 @@ const TERMINAL_SNAPSHOT_RESET: &[u8] = concat!(
|
|||||||
.as_bytes();
|
.as_bytes();
|
||||||
|
|
||||||
/// Seconds of silence from the server before the disconnect status line appears.
|
/// Seconds of silence from the server before the disconnect status line appears.
|
||||||
/// Short enough to give quick feedback on a lost link, long enough that a normal
|
/// Short enough to give quick feedback on a lost link, but later than the first
|
||||||
/// idle period (the run loop only pings after 2s of quiet) never flashes it.
|
/// idle ping. Using the same threshold as the first ping races its Pong: the bar
|
||||||
const DISCONNECT_STATUS_THRESHOLD_SECS: u64 = 2;
|
/// is painted just before the authenticated reply arrives, which then requests a
|
||||||
|
/// snapshot to restore row 1 and turns every healthy idle connection into a
|
||||||
|
/// periodic reconnect loop.
|
||||||
|
const DISCONNECT_STATUS_THRESHOLD_SECS: u64 = 3;
|
||||||
|
|
||||||
/// Mosh-style disconnect status bar with snapshot-backed restoration.
|
/// Mosh-style disconnect status bar with snapshot-backed restoration.
|
||||||
///
|
///
|
||||||
@@ -11161,8 +11188,8 @@ const TERMINAL_CLEANUP: &[u8] = concat!(
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{
|
use super::{
|
||||||
ALT_SCREEN_IDLE_REPAINT_AFTER, CachedCredential, DisconnectStatus, DynamicForward,
|
CachedCredential, DisconnectStatus, DynamicForward, FRAME_GAP_RESYNC_AFTER_MS, FrameBuffer,
|
||||||
FRAME_GAP_RESYNC_AFTER_MS, FrameBuffer, LOCAL_SLEEP_REPAINT_AFTER,
|
LOCAL_SLEEP_REPAINT_AFTER, LOCAL_SLEEP_REPAINT_RETRY_AFTER,
|
||||||
LOCAL_SLEEP_REPAINT_RETRY_WINDOW, LocalForward, MAX_PENDING_USER_INPUT_BYTES,
|
LOCAL_SLEEP_REPAINT_RETRY_WINDOW, LocalForward, MAX_PENDING_USER_INPUT_BYTES,
|
||||||
NativeIdentityContext, POST_RECONNECT_STALE_INPUT_GRACE, POST_SUBMIT_ALL_INPUT_HOLD,
|
NativeIdentityContext, POST_RECONNECT_STALE_INPUT_GRACE, POST_SUBMIT_ALL_INPUT_HOLD,
|
||||||
PendingStreamControl, PendingStreamOpen, PendingWindowAdjust, PredictMode, Predictor,
|
PendingStreamControl, PendingStreamOpen, PendingWindowAdjust, PredictMode, Predictor,
|
||||||
@@ -14433,86 +14460,35 @@ mod tests {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn idle_repaint_runs_for_stale_alternate_screen_or_sleep_gap() {
|
fn idle_repaint_runs_only_after_sleep_or_an_armed_wake_retry() {
|
||||||
let now = Instant::now();
|
let now = Instant::now();
|
||||||
let stale = now - ALT_SCREEN_IDLE_REPAINT_AFTER - Duration::from_secs(1);
|
let stale = now - LOCAL_SLEEP_REPAINT_RETRY_AFTER - Duration::from_secs(1);
|
||||||
let recent = now - Duration::from_secs(1);
|
|
||||||
let just_attempted = now - Duration::from_millis(500);
|
let just_attempted = now - Duration::from_millis(500);
|
||||||
assert!(should_repaint_idle_terminal(
|
assert!(!should_repaint_idle_terminal(
|
||||||
true,
|
|
||||||
stale,
|
|
||||||
stale,
|
stale,
|
||||||
Duration::from_secs(1),
|
Duration::from_secs(1),
|
||||||
None,
|
None,
|
||||||
now
|
now
|
||||||
));
|
));
|
||||||
assert!(should_repaint_idle_terminal(
|
assert!(should_repaint_idle_terminal(
|
||||||
false,
|
|
||||||
recent,
|
|
||||||
stale,
|
stale,
|
||||||
LOCAL_SLEEP_REPAINT_AFTER + Duration::from_secs(1),
|
LOCAL_SLEEP_REPAINT_AFTER + Duration::from_secs(1),
|
||||||
None,
|
None,
|
||||||
now
|
now
|
||||||
));
|
));
|
||||||
assert!(should_repaint_idle_terminal(
|
assert!(should_repaint_idle_terminal(
|
||||||
false,
|
stale,
|
||||||
recent,
|
Duration::from_secs(1),
|
||||||
recent,
|
|
||||||
LOCAL_SLEEP_REPAINT_AFTER + Duration::from_secs(1),
|
|
||||||
Some(now + LOCAL_SLEEP_REPAINT_RETRY_WINDOW),
|
Some(now + LOCAL_SLEEP_REPAINT_RETRY_WINDOW),
|
||||||
now
|
now
|
||||||
));
|
));
|
||||||
assert!(!should_repaint_idle_terminal(
|
assert!(!should_repaint_idle_terminal(
|
||||||
false,
|
|
||||||
recent,
|
|
||||||
just_attempted,
|
just_attempted,
|
||||||
LOCAL_SLEEP_REPAINT_AFTER + Duration::from_secs(1),
|
LOCAL_SLEEP_REPAINT_AFTER + Duration::from_secs(1),
|
||||||
Some(now + LOCAL_SLEEP_REPAINT_RETRY_WINDOW),
|
Some(now + LOCAL_SLEEP_REPAINT_RETRY_WINDOW),
|
||||||
now
|
now
|
||||||
));
|
));
|
||||||
assert!(!should_repaint_idle_terminal(
|
assert!(!should_repaint_idle_terminal(
|
||||||
false,
|
|
||||||
stale,
|
|
||||||
stale,
|
|
||||||
Duration::from_secs(1),
|
|
||||||
None,
|
|
||||||
now
|
|
||||||
));
|
|
||||||
assert!(!should_repaint_idle_terminal(
|
|
||||||
true,
|
|
||||||
recent,
|
|
||||||
stale,
|
|
||||||
Duration::from_secs(1),
|
|
||||||
None,
|
|
||||||
now
|
|
||||||
));
|
|
||||||
assert!(should_repaint_idle_terminal(
|
|
||||||
true,
|
|
||||||
stale,
|
|
||||||
recent,
|
|
||||||
LOCAL_SLEEP_REPAINT_AFTER + Duration::from_secs(1),
|
|
||||||
None,
|
|
||||||
now
|
|
||||||
));
|
|
||||||
assert!(should_repaint_idle_terminal(
|
|
||||||
false,
|
|
||||||
recent,
|
|
||||||
stale,
|
|
||||||
Duration::from_secs(1),
|
|
||||||
Some(now + LOCAL_SLEEP_REPAINT_RETRY_WINDOW),
|
|
||||||
now
|
|
||||||
));
|
|
||||||
assert!(!should_repaint_idle_terminal(
|
|
||||||
false,
|
|
||||||
recent,
|
|
||||||
just_attempted,
|
|
||||||
Duration::from_secs(1),
|
|
||||||
Some(now + LOCAL_SLEEP_REPAINT_RETRY_WINDOW),
|
|
||||||
now
|
|
||||||
));
|
|
||||||
assert!(!should_repaint_idle_terminal(
|
|
||||||
false,
|
|
||||||
recent,
|
|
||||||
stale,
|
stale,
|
||||||
Duration::from_secs(1),
|
Duration::from_secs(1),
|
||||||
Some(now - Duration::from_millis(1)),
|
Some(now - Duration::from_millis(1)),
|
||||||
|
|||||||
@@ -50,6 +50,9 @@ struct CachedCredentialWire {
|
|||||||
enum ServerObservation {
|
enum ServerObservation {
|
||||||
BulkSent,
|
BulkSent,
|
||||||
Input(Vec<u8>),
|
Input(Vec<u8>),
|
||||||
|
Ping,
|
||||||
|
Reconnected,
|
||||||
|
RenderResynced,
|
||||||
Resize(u16, u16),
|
Resize(u16, u16),
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -66,7 +69,7 @@ fn native_client_terminal_round_trip_is_platform_complete() {
|
|||||||
.set_read_timeout(Some(Duration::from_secs(8)))
|
.set_read_timeout(Some(Duration::from_secs(8)))
|
||||||
.unwrap();
|
.unwrap();
|
||||||
let port = socket.local_addr().unwrap().port();
|
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 (observation_tx, observation_rx) = mpsc::channel();
|
||||||
let server = thread::spawn(move || run_fake_terminal_server(socket, observation_tx));
|
let server = thread::spawn(move || run_fake_terminal_server(socket, observation_tx));
|
||||||
@@ -175,7 +178,7 @@ fn terminal_output_backpressure_does_not_block_input_transport() {
|
|||||||
.set_read_timeout(Some(Duration::from_secs(8)))
|
.set_read_timeout(Some(Duration::from_secs(8)))
|
||||||
.unwrap();
|
.unwrap();
|
||||||
let port = socket.local_addr().unwrap().port();
|
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 (observation_tx, observation_rx) = mpsc::channel();
|
||||||
let server = thread::spawn(move || run_backpressure_server(socket, observation_tx, PROBE));
|
let server = thread::spawn(move || run_backpressure_server(socket, observation_tx, PROBE));
|
||||||
|
|
||||||
@@ -224,12 +227,242 @@ fn terminal_output_backpressure_does_not_block_input_transport() {
|
|||||||
assert!(status.success(), "Dosh client exited with {status:?}");
|
assert!(status.success(), "Dosh client exited with {status:?}");
|
||||||
}
|
}
|
||||||
|
|
||||||
fn write_client_fixture(home: &Path, cache: &Path, port: u16) {
|
#[test]
|
||||||
|
fn idle_reconnect_restores_snapshot_and_orders_reordered_frames() {
|
||||||
|
const PROBE: &[u8] = b"DOSH_RECONNECT_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(8)))
|
||||||
|
.unwrap();
|
||||||
|
let port = socket.local_addr().unwrap().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));
|
||||||
|
|
||||||
|
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);
|
||||||
|
|
||||||
|
let output = Arc::new(Mutex::new(Vec::new()));
|
||||||
|
let reader_output = Arc::clone(&output);
|
||||||
|
let reader_thread = thread::spawn(move || {
|
||||||
|
let mut buf = [0u8; 4096];
|
||||||
|
loop {
|
||||||
|
match reader.read(&mut buf) {
|
||||||
|
Ok(0) | Err(_) => break,
|
||||||
|
Ok(n) => reader_output.lock().unwrap().extend_from_slice(&buf[..n]),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
wait_for_output(&output, b"DOSH_RECONNECT_READY", Duration::from_secs(3));
|
||||||
|
wait_for_observation(&observation_rx, Duration::from_secs(4), |observation| {
|
||||||
|
matches!(observation, ServerObservation::Reconnected)
|
||||||
|
});
|
||||||
|
wait_for_output(&output, b"DOSH_RECONNECT_SNAPSHOT", Duration::from_secs(3));
|
||||||
|
|
||||||
|
writer.write_all(PROBE).unwrap();
|
||||||
|
writer.flush().unwrap();
|
||||||
|
wait_for_input(&observation_rx, PROBE, Duration::from_secs(2));
|
||||||
|
wait_for_output(
|
||||||
|
&output,
|
||||||
|
b"DOSH_ORDER_FIRST_DOSH_ORDER_SECOND_DOSH_RECONNECT_DONE",
|
||||||
|
Duration::from_secs(3),
|
||||||
|
);
|
||||||
|
|
||||||
|
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:?}");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn authenticated_idle_pongs_do_not_trigger_snapshot_reconnects() {
|
||||||
|
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_millis(100)))
|
||||||
|
.unwrap();
|
||||||
|
let port = socket.local_addr().unwrap().port();
|
||||||
|
write_client_fixture(&home, &cache, port, 5);
|
||||||
|
let config_path = home.join(".config/dosh/client.toml");
|
||||||
|
let mut config: ClientConfig =
|
||||||
|
toml::from_str(&fs::read_to_string(&config_path).unwrap()).unwrap();
|
||||||
|
config.disconnect_status = true;
|
||||||
|
fs::write(&config_path, toml::to_string(&config).unwrap()).unwrap();
|
||||||
|
|
||||||
|
let (observation_tx, observation_rx) = mpsc::channel();
|
||||||
|
let server = thread::spawn(move || run_idle_keepalive_server(socket, observation_tx));
|
||||||
|
|
||||||
|
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 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);
|
||||||
|
|
||||||
|
let output = Arc::new(Mutex::new(Vec::new()));
|
||||||
|
let reader_output = Arc::clone(&output);
|
||||||
|
let reader_thread = thread::spawn(move || {
|
||||||
|
let mut buf = [0u8; 4096];
|
||||||
|
loop {
|
||||||
|
match reader.read(&mut buf) {
|
||||||
|
Ok(0) | Err(_) => break,
|
||||||
|
Ok(n) => reader_output.lock().unwrap().extend_from_slice(&buf[..n]),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
let status = child.wait().unwrap();
|
||||||
|
drop(pair.master);
|
||||||
|
reader_thread.join().unwrap();
|
||||||
|
server.join().unwrap();
|
||||||
|
assert!(status.success(), "Dosh client exited with {status:?}");
|
||||||
|
|
||||||
|
let observations: Vec<_> = observation_rx.try_iter().collect();
|
||||||
|
assert!(
|
||||||
|
observations
|
||||||
|
.iter()
|
||||||
|
.filter(|event| matches!(event, ServerObservation::Ping))
|
||||||
|
.count()
|
||||||
|
>= 2,
|
||||||
|
"idle session did not exercise repeated authenticated keepalives: {observations:?}"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
!observations
|
||||||
|
.iter()
|
||||||
|
.any(|event| matches!(event, ServerObservation::Reconnected)),
|
||||||
|
"healthy idle pongs caused a snapshot reconnect: {observations:?}"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
!contains(&output.lock().unwrap(), b"[dosh] reconnecting"),
|
||||||
|
"healthy idle session flashed the disconnect overlay"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[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 {
|
let config = ClientConfig {
|
||||||
server: "local".to_string(),
|
server: "local".to_string(),
|
||||||
dosh_host: Some("127.0.0.1".to_string()),
|
dosh_host: Some("127.0.0.1".to_string()),
|
||||||
dosh_port: port,
|
dosh_port: port,
|
||||||
cache_attach_tickets: false,
|
cache_attach_tickets: false,
|
||||||
|
reconnect_timeout_secs,
|
||||||
credential_cache: cache.to_string_lossy().to_string(),
|
credential_cache: cache.to_string_lossy().to_string(),
|
||||||
auth_preference: "ssh".to_string(),
|
auth_preference: "ssh".to_string(),
|
||||||
predict: false,
|
predict: false,
|
||||||
@@ -479,6 +712,265 @@ fn run_backpressure_server(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn run_reconnect_server(
|
||||||
|
socket: UdpSocket,
|
||||||
|
observations: mpsc::Sender<ServerObservation>,
|
||||||
|
expected_input: &[u8],
|
||||||
|
) {
|
||||||
|
let mut resume_count = 0u64;
|
||||||
|
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 => {
|
||||||
|
let plain =
|
||||||
|
protocol::decrypt_body(&packet, &SESSION_KEY, CLIENT_TO_SERVER).unwrap();
|
||||||
|
let request: protocol::ResumeRequest = protocol::from_body(&plain).unwrap();
|
||||||
|
assert_eq!(request.session, SESSION);
|
||||||
|
resume_count += 1;
|
||||||
|
let bytes = if resume_count == 1 {
|
||||||
|
b"DOSH_RECONNECT_READY".as_slice()
|
||||||
|
} else {
|
||||||
|
observations.send(ServerObservation::Reconnected).unwrap();
|
||||||
|
b"DOSH_RECONNECT_SNAPSHOT".as_slice()
|
||||||
|
};
|
||||||
|
send_frame(
|
||||||
|
&socket,
|
||||||
|
source,
|
||||||
|
PacketKind::ResumeOk,
|
||||||
|
resume_count,
|
||||||
|
10,
|
||||||
|
bytes,
|
||||||
|
true,
|
||||||
|
false,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
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 contains(&received_input, expected_input) {
|
||||||
|
send_frame(
|
||||||
|
&socket,
|
||||||
|
source,
|
||||||
|
PacketKind::Frame,
|
||||||
|
4,
|
||||||
|
12,
|
||||||
|
b"DOSH_ORDER_SECOND_",
|
||||||
|
false,
|
||||||
|
false,
|
||||||
|
);
|
||||||
|
send_frame(
|
||||||
|
&socket,
|
||||||
|
source,
|
||||||
|
PacketKind::Frame,
|
||||||
|
3,
|
||||||
|
11,
|
||||||
|
b"DOSH_ORDER_FIRST_",
|
||||||
|
false,
|
||||||
|
false,
|
||||||
|
);
|
||||||
|
send_frame(
|
||||||
|
&socket,
|
||||||
|
source,
|
||||||
|
PacketKind::Frame,
|
||||||
|
5,
|
||||||
|
13,
|
||||||
|
b"DOSH_RECONNECT_DONE",
|
||||||
|
false,
|
||||||
|
true,
|
||||||
|
);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
PacketKind::Ping | PacketKind::Ack => {}
|
||||||
|
_ => {}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn run_idle_keepalive_server(socket: UdpSocket, observations: mpsc::Sender<ServerObservation>) {
|
||||||
|
let started = Instant::now();
|
||||||
|
let mut peer = None;
|
||||||
|
let mut attached = false;
|
||||||
|
let mut server_seq = 1u64;
|
||||||
|
let mut buf = [0u8; 65535];
|
||||||
|
loop {
|
||||||
|
if started.elapsed() >= Duration::from_secs(7)
|
||||||
|
&& let Some(source) = peer
|
||||||
|
{
|
||||||
|
server_seq += 1;
|
||||||
|
send_frame(
|
||||||
|
&socket,
|
||||||
|
source,
|
||||||
|
PacketKind::Frame,
|
||||||
|
server_seq,
|
||||||
|
11,
|
||||||
|
b"DOSH_IDLE_KEEPALIVE_DONE",
|
||||||
|
false,
|
||||||
|
true,
|
||||||
|
);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
|
let (n, source) = match socket.recv_from(&mut buf) {
|
||||||
|
Ok(value) => value,
|
||||||
|
Err(err)
|
||||||
|
if matches!(
|
||||||
|
err.kind(),
|
||||||
|
std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
|
||||||
|
) =>
|
||||||
|
{
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
Err(err) => panic!("idle keepalive server receive failed: {err}"),
|
||||||
|
};
|
||||||
|
let packet = protocol::decode(&buf[..n]).unwrap();
|
||||||
|
match packet.header.kind {
|
||||||
|
PacketKind::ResumeRequest => {
|
||||||
|
let plain =
|
||||||
|
protocol::decrypt_body(&packet, &SESSION_KEY, CLIENT_TO_SERVER).unwrap();
|
||||||
|
let request: protocol::ResumeRequest = protocol::from_body(&plain).unwrap();
|
||||||
|
assert_eq!(request.session, SESSION);
|
||||||
|
if attached {
|
||||||
|
observations.send(ServerObservation::Reconnected).unwrap();
|
||||||
|
}
|
||||||
|
attached = true;
|
||||||
|
peer = Some(source);
|
||||||
|
send_frame(
|
||||||
|
&socket,
|
||||||
|
source,
|
||||||
|
PacketKind::ResumeOk,
|
||||||
|
server_seq,
|
||||||
|
10,
|
||||||
|
b"DOSH_IDLE_KEEPALIVE_READY",
|
||||||
|
true,
|
||||||
|
false,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
PacketKind::Ping => {
|
||||||
|
protocol::decrypt_body(&packet, &SESSION_KEY, CLIENT_TO_SERVER).unwrap();
|
||||||
|
observations.send(ServerObservation::Ping).unwrap();
|
||||||
|
server_seq += 1;
|
||||||
|
let pong = protocol::encode_encrypted(
|
||||||
|
PacketKind::Pong,
|
||||||
|
CLIENT_ID,
|
||||||
|
server_seq,
|
||||||
|
packet.header.seq,
|
||||||
|
&SESSION_KEY,
|
||||||
|
SERVER_TO_CLIENT,
|
||||||
|
b"",
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
socket.send_to(&pong, source).unwrap();
|
||||||
|
}
|
||||||
|
PacketKind::Ack => {}
|
||||||
|
_ => {}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
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)]
|
#[allow(clippy::too_many_arguments)]
|
||||||
fn send_frame(
|
fn send_frame(
|
||||||
socket: &UdpSocket,
|
socket: &UdpSocket,
|
||||||
|
|||||||
Reference in New Issue
Block a user