pub mod cli;
mod render;
pub use cli::{connect, run_id, ConnectArgs, IdArgs};
use std::io::Write;
use std::time::Duration;
use crate::input::UserInput;
use crate::predict::{DisplayPreference, Overlay, PredictionEngine};
use crate::ssp::{RecvOutcome, Transport, SHUTDOWN_SENTINEL};
use crate::terminal::TerminalScreen;
use crate::transport_iroh::{IrohChannel, MonoClock, ALPN};
use anyhow::Context;
use iroh::{Endpoint, EndpointAddr};
use termina::escape::csi::{Csi, DecPrivateMode, DecPrivateModeCode, Mode};
use termina::{PlatformTerminal, Terminal as _};
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
pub use render::WindowState;
const KOH_TITLE_PREFIX: &str = "[koh] ";
pub(crate) const ESCAPE_PREFIX: u8 = 0x1e;
pub(crate) const SUSPEND_KEY: u8 = 0x1a;
const RESET_FORWARDED_MODES: &[u8] =
b"\x1b[?9l\x1b[?2004l\x1b[?1000l\x1b[?1002l\x1b[?1003l\x1b[?1005l\x1b[?1006l\x1b[?1l\x1b>";
const RECONNECT_CONNECT_TIMEOUT: Duration = Duration::from_secs(15);
const RECONNECT_BACKOFF_BASE_MS: u64 = 500;
const RECONNECT_BACKOFF_MAX_MS: u64 = 8_000;
const MIN_CONNECTION_DWELL_MS: u64 = 5_000;
const LINK_DOWN_GRACE_MS: u64 = crate::ssp::ACK_INTERVAL * 3;
const STALE_AFTER_FREEZE: Duration = Duration::from_secs(20);
fn looks_like_resume_from_freeze(wall_gap: Duration) -> bool {
wall_gap >= STALE_AFTER_FREEZE
}
pub struct IrohConnector {
endpoint: Endpoint,
target: EndpointAddr,
}
impl IrohConnector {
pub fn new(endpoint: Endpoint, target: EndpointAddr) -> Self {
Self { endpoint, target }
}
pub async fn connect(&self) -> anyhow::Result<IrohChannel> {
let conn = self
.endpoint
.connect(self.target.clone(), ALPN)
.await
.context("connecting to server (is your id on its allowlist?)")?;
if let Err(e) = crate::transport_iroh::admission::await_admission(&conn).await {
return Err(match server_close_reason(&conn) {
Some(reason) => anyhow::Error::new(e)
.context(format!("server rejected the connection: {reason}")),
None => anyhow::Error::new(e)
.context("server did not admit the connection (is your id on its allowlist?)"),
});
}
Ok(IrohChannel::new(conn))
}
}
fn server_close_reason(conn: &iroh::endpoint::Connection) -> Option<String> {
use iroh::endpoint::{ApplicationClose, ConnectionError};
let ConnectionError::ApplicationClosed(ApplicationClose { reason, .. }) =
conn.close_reason()?
else {
return None;
};
let cleaned: String = String::from_utf8_lossy(&reason)
.chars()
.filter(|c| !c.is_control())
.take(80)
.collect();
(!cleaned.is_empty()).then_some(cleaned)
}
fn backoff_ms(attempt: u32) -> u64 {
(RECONNECT_BACKOFF_BASE_MS << attempt.min(4)).min(RECONNECT_BACKOFF_MAX_MS)
}
const fn next_attempt_after_drop(attempt: u32, dwell_ms: u64) -> u32 {
if dwell_ms >= MIN_CONNECTION_DWELL_MS {
0
} else {
attempt.saturating_add(1)
}
}
fn escape_quit(chunk: &[u8], pending: &mut bool) -> bool {
for &b in chunk {
if *pending {
*pending = false;
if b == b'.' {
return true;
}
} else if b == ESCAPE_PREFIX {
*pending = true;
}
}
false
}
pub trait ClientTerminal {
fn render(
&mut self,
screen: &vt100::Screen,
overlay: &Overlay,
status: Option<&str>,
win: render::WindowState<'_>,
) -> std::io::Result<()>;
fn size(&self) -> std::io::Result<(u16, u16)>;
fn suspend_resume(&mut self) -> std::io::Result<()> {
Ok(())
}
}
pub struct TerminaTerminal {
term: PlatformTerminal,
oob: render::OutOfBand,
}
impl TerminaTerminal {
pub fn enter(clipboard_enabled: bool) -> std::io::Result<Self> {
let mut term = PlatformTerminal::new()?;
term.enter_raw_mode()?;
let mut this = Self {
term,
oob: render::OutOfBand::with_title_prefix(KOH_TITLE_PREFIX.to_string())
.with_clipboard(clipboard_enabled),
};
this.enter_screen()?;
Ok(this)
}
fn enter_screen(&mut self) -> std::io::Result<()> {
write!(
self.term,
"{}{}",
Csi::Mode(Mode::SetDecPrivateMode(DecPrivateMode::Code(
DecPrivateModeCode::ClearAndEnableAlternateScreen
))),
Csi::Mode(Mode::ResetDecPrivateMode(DecPrivateMode::Code(
DecPrivateModeCode::ShowCursor
))),
)?;
self.term.flush()
}
fn leave_screen(&mut self) -> std::io::Result<()> {
self.term.write_all(RESET_FORWARDED_MODES)?;
write!(
self.term,
"{}{}",
Csi::Mode(Mode::SetDecPrivateMode(DecPrivateMode::Code(
DecPrivateModeCode::ShowCursor
))),
Csi::Mode(Mode::ResetDecPrivateMode(DecPrivateMode::Code(
DecPrivateModeCode::ClearAndEnableAlternateScreen
))),
)?;
self.term.flush()
}
}
impl ClientTerminal for TerminaTerminal {
fn render(
&mut self,
screen: &vt100::Screen,
overlay: &Overlay,
status: Option<&str>,
win: render::WindowState<'_>,
) -> std::io::Result<()> {
self.oob.emit(&mut self.term, screen, win)?;
render::render(&mut self.term, screen, overlay, status)
}
fn size(&self) -> std::io::Result<(u16, u16)> {
let d = self.term.get_dimensions()?;
Ok((d.rows, d.cols))
}
fn suspend_resume(&mut self) -> std::io::Result<()> {
self.leave_screen()?;
self.term.enter_cooked_mode()?;
let _ = writeln!(self.term, "\n[koh suspended — run `fg` to resume]");
let _ = self.term.flush();
nix::sys::signal::raise(nix::sys::signal::Signal::SIGTSTP)
.map_err(std::io::Error::other)?;
self.term.enter_raw_mode()?;
self.enter_screen()?;
self.oob.invalidate();
Ok(())
}
}
impl Drop for TerminaTerminal {
fn drop(&mut self) {
let _ = self.leave_screen();
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum InputOutcome {
Quit,
Suspend,
Forwarded,
}
#[derive(Debug, Default)]
pub struct TickResult {
pub outgoing: Vec<Vec<u8>>,
pub wait_ms: u64,
pub status: Option<String>,
pub ended: Option<Option<u32>>,
}
pub struct ClientSession {
transport: Transport<UserInput, TerminalScreen>,
predictor: PredictionEngine,
pending_escape: bool,
dirty: bool,
status_was_shown: bool,
}
impl ClientSession {
pub fn new(
now: u64,
mtu: usize,
pref: DisplayPreference,
initial_rows: u16,
initial_cols: u16,
) -> Self {
let mut transport = Transport::<UserInput, TerminalScreen>::new(now, mtu);
transport.set_connected(true);
transport
.current_mut()
.push_resize(initial_rows, initial_cols);
let predictor = PredictionEngine::new(pref);
Self {
transport,
predictor,
pending_escape: false,
dirty: true,
status_was_shown: false,
}
}
pub fn on_input(&mut self, now: u64, bytes: &[u8]) -> InputOutcome {
let mut quit = false;
let mut suspend = false;
let mut fwd: Vec<u8> = Vec::with_capacity(bytes.len());
for &b in bytes {
if self.pending_escape {
self.pending_escape = false;
if b == b'.' {
quit = true;
break;
}
if b == SUSPEND_KEY {
suspend = true;
break;
}
fwd.push(ESCAPE_PREFIX);
fwd.push(b);
} else if b == ESCAPE_PREFIX {
self.pending_escape = true;
} else {
fwd.push(b);
}
}
if quit {
return InputOutcome::Quit;
}
if !fwd.is_empty() {
self.predictor
.set_local_frame_sent(self.transport.newest_sent_num());
self.predictor
.set_srtt(self.transport.send_interval() as f64);
let screen = self.transport.remote_state().screen();
for &b in &fwd {
self.predictor.new_user_byte(now, b, screen);
}
self.transport.current_mut().push_bytes(&fwd);
self.dirty = true;
}
if suspend {
return InputOutcome::Suspend;
}
InputOutcome::Forwarded
}
pub fn on_datagram(&mut self, now: u64, bytes: &[u8]) {
if self.transport.recv(now, bytes) == RecvOutcome::NewState {
let echo_ack = self.transport.remote_state().echo_ack();
self.predictor.set_local_frame_late_acked(echo_ack);
self.predictor
.set_srtt(self.transport.send_interval() as f64);
let screen = self.transport.remote_state().screen();
self.predictor.cull(now, screen);
self.dirty = true;
}
}
pub fn on_resize(&mut self, rows: u16, cols: u16) {
self.transport.current_mut().push_resize(rows, cols);
self.predictor.reset();
self.dirty = true;
}
pub fn on_tick(&mut self, now: u64, mtu: usize, rtt_ms: Option<f64>) -> TickResult {
self.transport.set_mtu(mtu);
if let Some(rtt) = rtt_ms {
self.transport.observe_rtt(rtt);
}
if self
.predictor
.tick(now, self.transport.remote_state().screen())
{
self.dirty = true;
}
let outgoing = self.transport.tick(now);
let status = if self.transport.last_heard() > 0
&& !self.transport.link_up_within(now, LINK_DOWN_GRACE_MS)
{
let since = now.saturating_sub(self.transport.last_heard());
Some(format!("[koh] link down — resuming… {}s", since / 1000))
} else {
None
};
let ended = (self.transport.remote_num() == SHUTDOWN_SENTINEL)
.then(|| self.transport.remote_state().exit_code());
let wait_ms = self.transport.wait_time(now).min(50);
TickResult {
outgoing,
wait_ms,
status,
ended,
}
}
pub fn screen(&self) -> &vt100::Screen {
self.transport.remote_state().screen()
}
pub fn overlay(&self) -> Overlay {
self.predictor
.overlay(self.transport.remote_state().screen())
}
pub fn window_state(&self) -> render::WindowState<'_> {
let ts = self.transport.remote_state();
render::WindowState {
title: ts.title(),
icon: ts.icon(),
clipboard: ts.clipboard(),
bell_count: ts.bell_count(),
}
}
}
#[expect(
clippy::too_many_arguments,
reason = "the I/O shell wires up the channel, connector, prediction policy, size, the two \
input/resize channels, the terminal, and the shutdown token — each a distinct \
collaborator; bundling them into a struct would only move the list, not shorten it"
)]
pub async fn run_client<T: ClientTerminal>(
initial: IrohChannel,
connector: IrohConnector,
pref: DisplayPreference,
initial_size: (u16, u16),
mut input_rx: mpsc::Receiver<Vec<u8>>,
mut resize_rx: mpsc::Receiver<()>,
mut term: T,
shutdown: CancellationToken,
) -> anyhow::Result<Option<u32>> {
let clock = MonoClock::new();
let mut channel = initial;
let mut attempt: u32 = 0;
loop {
let (rows, cols) = term.size().unwrap_or(initial_size);
let mut session = ClientSession::new(
clock.now_ms(),
channel.max_datagram_size(),
pref,
rows,
cols,
);
let conn_started = clock.now_ms();
match drive_connection(
&channel,
&mut session,
&mut term,
&mut input_rx,
&mut resize_rx,
&clock,
&shutdown,
)
.await?
{
Disposition::Quit => {
channel.close(0, b"client exit");
return Ok(None);
}
Disposition::Ended(code) => {
channel.close(0, b"client exit");
return Ok(code);
}
Disposition::LinkLost => {
channel.close(0, b"reconnecting");
let dwell = clock.now_ms().saturating_sub(conn_started);
attempt = next_attempt_after_drop(attempt, dwell);
match reconnect(
&connector,
&mut term,
&mut input_rx,
&session,
&clock,
&shutdown,
&mut attempt,
)
.await
{
ReconnectOutcome::Connected(c) => channel = c,
ReconnectOutcome::Quit => return Ok(None),
}
}
}
}
}
enum Disposition {
Quit,
Ended(Option<u32>),
LinkLost,
}
async fn drive_connection<T: ClientTerminal>(
channel: &IrohChannel,
session: &mut ClientSession,
term: &mut T,
input_rx: &mut mpsc::Receiver<Vec<u8>>,
resize_rx: &mut mpsc::Receiver<()>,
clock: &MonoClock,
shutdown: &CancellationToken,
) -> anyhow::Result<Disposition> {
let mut last_wall = std::time::SystemTime::now();
let mut last_logged_rtt: Option<f64> = None;
loop {
let wall_now = std::time::SystemTime::now();
let wall_gap = wall_now.duration_since(last_wall).unwrap_or(Duration::ZERO);
last_wall = wall_now;
if looks_like_resume_from_freeze(wall_gap) {
tracing::info!(
frozen_secs = wall_gap.as_secs(),
"detected resume from a process freeze (suspend/screen-off); forcing a reconnect"
);
return Ok(Disposition::LinkLost);
}
let now = clock.now_ms();
let rtt = channel.rtt_ms();
if let Some(ms) = rtt {
if last_logged_rtt.is_none_or(|prev| (prev - ms).abs() >= 30.0) {
tracing::debug!(rtt_ms = ms, "link rtt");
last_logged_rtt = Some(ms);
}
}
let tick = session.on_tick(now, channel.max_datagram_size(), rtt);
for datagram in &tick.outgoing {
channel.send(datagram);
}
let status_now = tick.status.is_some();
if session.dirty || status_now || session.status_was_shown {
term.render(
session.screen(),
&session.overlay(),
tick.status.as_deref(),
session.window_state(),
)?;
session.status_was_shown = status_now;
session.dirty = false;
}
if let Some(code) = tick.ended {
let _ = term.render(
session.screen(),
&Overlay::empty(),
Some("[koh] session ended"),
session.window_state(),
);
tokio::select! {
() = tokio::time::sleep(Duration::from_millis(400)) => {}
() = shutdown.cancelled() => {}
}
return Ok(Disposition::Ended(code));
}
tokio::select! {
biased;
maybe = input_rx.recv() => {
match maybe {
Some(chunk) => match session.on_input(clock.now_ms(), &chunk) {
InputOutcome::Quit => return Ok(Disposition::Quit),
InputOutcome::Suspend => {
term.suspend_resume()?;
session.dirty = true;
last_wall = std::time::SystemTime::now();
}
InputOutcome::Forwarded => {}
},
None => return Ok(Disposition::Quit), }
}
_ = shutdown.cancelled() => return Ok(Disposition::Quit),
dg = channel.recv() => {
match dg {
Ok(bytes) => session.on_datagram(clock.now_ms(), &bytes),
Err(e) => {
tracing::info!(reason = %e, "link lost; will reconnect");
return Ok(Disposition::LinkLost);
}
}
}
maybe = resize_rx.recv() => {
if maybe.is_some() {
if let Ok((rows, cols)) = term.size() {
session.on_resize(rows, cols);
}
}
}
_ = tokio::time::sleep(Duration::from_millis(tick.wait_ms)) => {}
}
}
}
enum ReconnectOutcome {
Connected(IrohChannel),
Quit,
}
async fn reconnect<T: ClientTerminal>(
connector: &IrohConnector,
term: &mut T,
input_rx: &mut mpsc::Receiver<Vec<u8>>,
last: &ClientSession,
clock: &MonoClock,
shutdown: &CancellationToken,
attempt: &mut u32,
) -> ReconnectOutcome {
let started = clock.now_ms();
let mut pending_escape = false;
'attempt: loop {
if *attempt > 0 {
let wait_until = clock.now_ms().saturating_add(backoff_ms(*attempt));
while clock.now_ms() < wait_until {
let secs = clock.now_ms().saturating_sub(started) / 1000;
let banner =
format!("[koh] disconnected — reconnecting… {secs}s (Ctrl-^ . to quit)");
let _ = term.render(
last.screen(),
&Overlay::empty(),
Some(banner.as_str()),
last.window_state(),
);
let remaining = wait_until.saturating_sub(clock.now_ms());
tokio::select! {
biased;
maybe = input_rx.recv() => match maybe {
Some(chunk) => {
if escape_quit(&chunk, &mut pending_escape) {
return ReconnectOutcome::Quit;
}
}
None => return ReconnectOutcome::Quit,
},
_ = shutdown.cancelled() => return ReconnectOutcome::Quit,
_ = tokio::time::sleep(Duration::from_millis(remaining.min(1000))) => {}
}
}
}
let dial = tokio::time::timeout(RECONNECT_CONNECT_TIMEOUT, connector.connect());
tokio::pin!(dial);
loop {
let secs = clock.now_ms().saturating_sub(started) / 1000;
let banner = format!("[koh] disconnected — reconnecting… {secs}s (Ctrl-^ . to quit)");
let _ = term.render(
last.screen(),
&Overlay::empty(),
Some(banner.as_str()),
last.window_state(),
);
tokio::select! {
biased;
maybe = input_rx.recv() => {
match maybe {
Some(chunk) => {
if escape_quit(&chunk, &mut pending_escape) {
return ReconnectOutcome::Quit;
}
}
None => return ReconnectOutcome::Quit, }
}
res = &mut dial => {
match res {
Ok(Ok(channel)) => return ReconnectOutcome::Connected(channel),
Ok(Err(e)) => tracing::info!(reason = %e, attempt = *attempt, "reconnect dial failed"),
Err(_) => tracing::info!(attempt = *attempt, "reconnect dial timed out"),
}
*attempt = (*attempt).saturating_add(1);
continue 'attempt;
}
_ = shutdown.cancelled() => return ReconnectOutcome::Quit,
_ = tokio::time::sleep(Duration::from_secs(1)) => { }
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::input::InputEvent;
use crate::terminal::ServerTerminal;
fn drive_until_nonempty(t: &mut Transport<TerminalScreen, UserInput>) -> Vec<Vec<u8>> {
let mut now = 0u64;
loop {
now += 25;
let out = t.tick(now);
if !out.is_empty() || now > 5_000 {
return out;
}
}
}
fn new_session() -> ClientSession {
ClientSession::new(0, 1200, DisplayPreference::Always, 24, 80)
}
#[test]
fn escape_prefix_dot_quits_and_plain_bytes_forward() {
let mut s = new_session();
assert_eq!(s.on_input(0, b"ls\r"), InputOutcome::Forwarded);
let typed: Vec<u8> = s
.transport
.current()
.events()
.iter()
.filter_map(|e| match e {
InputEvent::Byte(b) => Some(*b),
InputEvent::Resize { .. } => None,
})
.collect();
assert_eq!(
typed, b"ls\r",
"forwarded bytes land in transport.current()"
);
assert_eq!(s.on_input(0, &[ESCAPE_PREFIX, b'.']), InputOutcome::Quit);
}
#[test]
fn escape_prefix_ctrl_z_suspends() {
let mut s = new_session();
assert_eq!(
s.on_input(0, &[ESCAPE_PREFIX, SUSPEND_KEY]),
InputOutcome::Suspend
);
assert_eq!(s.on_input(0, &[ESCAPE_PREFIX]), InputOutcome::Forwarded);
assert_eq!(s.on_input(0, &[SUSPEND_KEY]), InputOutcome::Suspend);
}
#[test]
fn bytes_before_suspend_escape_are_forwarded_first() {
let mut s = new_session();
assert_eq!(
s.on_input(0, &[b'h', b'i', ESCAPE_PREFIX, SUSPEND_KEY]),
InputOutcome::Suspend
);
let typed: Vec<u8> = s
.transport
.current()
.events()
.iter()
.filter_map(|e| match e {
InputEvent::Byte(b) => Some(*b),
InputEvent::Resize { .. } => None,
})
.collect();
assert_eq!(
typed, b"hi",
"pre-escape bytes are forwarded before suspending"
);
}
#[test]
fn escape_quit_matches_the_session_machine_across_chunks() {
let mut p = false;
assert!(!escape_quit(b"hello", &mut p), "plain bytes never quit");
assert!(!p);
assert!(escape_quit(&[ESCAPE_PREFIX, b'.'], &mut p));
p = false;
assert!(!escape_quit(&[ESCAPE_PREFIX], &mut p));
assert!(p, "a lone prefix leaves us pending");
assert!(escape_quit(b".", &mut p));
p = false;
assert!(!escape_quit(&[ESCAPE_PREFIX, b'x'], &mut p));
assert!(!p, "prefix + non-dot resets pending");
assert!(!escape_quit(b".", &mut p), "a later lone '.' must not quit");
}
#[test]
fn reconnect_backoff_grows_then_caps() {
assert_eq!(backoff_ms(1), 1_000);
assert_eq!(backoff_ms(2), 2_000);
assert_eq!(backoff_ms(3), 4_000);
assert_eq!(backoff_ms(4), RECONNECT_BACKOFF_MAX_MS);
assert_eq!(backoff_ms(5), RECONNECT_BACKOFF_MAX_MS);
assert_eq!(
backoff_ms(99),
RECONNECT_BACKOFF_MAX_MS,
"shift is clamped, no overflow"
);
}
#[test]
fn dwell_gate_resets_on_proven_connection_and_climbs_on_flap() {
assert_eq!(next_attempt_after_drop(0, MIN_CONNECTION_DWELL_MS), 0);
assert_eq!(next_attempt_after_drop(5, MIN_CONNECTION_DWELL_MS), 0);
assert_eq!(
next_attempt_after_drop(5, MIN_CONNECTION_DWELL_MS + 10_000),
0
);
assert_eq!(next_attempt_after_drop(0, 0), 1);
assert_eq!(next_attempt_after_drop(3, MIN_CONNECTION_DWELL_MS - 1), 4);
assert_eq!(next_attempt_after_drop(u32::MAX, 0), u32::MAX);
}
#[test]
fn freeze_detection_fires_only_on_a_real_suspend_gap() {
assert!(!looks_like_resume_from_freeze(Duration::from_millis(0)));
assert!(!looks_like_resume_from_freeze(Duration::from_millis(50)));
assert!(!looks_like_resume_from_freeze(Duration::from_secs(5)));
assert_eq!(STALE_AFTER_FREEZE, Duration::from_secs(20));
assert!(!looks_like_resume_from_freeze(Duration::from_secs(19)));
assert!(looks_like_resume_from_freeze(STALE_AFTER_FREEZE));
assert!(looks_like_resume_from_freeze(Duration::from_secs(300)));
}
#[test]
fn lone_escape_prefix_then_other_byte_forwards_both() {
let mut s = new_session();
assert_eq!(s.on_input(0, &[ESCAPE_PREFIX]), InputOutcome::Forwarded);
assert_eq!(s.on_input(0, b"x"), InputOutcome::Forwarded);
let typed: Vec<u8> = s
.transport
.current()
.events()
.iter()
.filter_map(|e| match e {
InputEvent::Byte(b) => Some(*b),
InputEvent::Resize { .. } => None,
})
.collect();
assert_eq!(
typed,
[ESCAPE_PREFIX, b'x'],
"escaped non-dot byte passes through literally"
);
}
#[test]
fn on_datagram_new_state_marks_dirty_and_culls_predictor() {
let mut s = new_session();
s.on_input(0, b"x");
s.dirty = false; assert!(
s.overlay().is_empty(),
"the first keystroke stays hidden until confirmed"
);
assert_eq!(
s.predictor.confirmed_epoch(),
0,
"nothing is confirmed before the server frame arrives"
);
let mut emu = ServerTerminal::new(24, 80, 0);
emu.process(b"x");
emu.register_input_frame(1, 0);
emu.set_echo_ack(100);
let mut server = Transport::<TerminalScreen, UserInput>::new(0, 1200);
server.set_connected(true);
server.observe_rtt(20.0);
*server.current_mut() = emu.snapshot();
for dg in drive_until_nonempty(&mut server) {
s.on_datagram(100, &dg);
}
assert!(
s.dirty,
"a new remote state must mark the client dirty (needs repaint)"
);
assert!(
s.screen().contents().contains('x'),
"the new state is applied to the screen"
);
assert_eq!(
s.predictor.confirmed_epoch(),
1,
"on_datagram's cull must grade the echoed 'x' Correct and advance the confirmed epoch"
);
s.on_input(110, b"y");
assert_eq!(
s.overlay().cell(0, 1).map(|c| c.glyph.as_str()),
Some("y"),
"typing after the confirmed echo is visible (the prior prediction was culled)"
);
}
#[test]
fn on_tick_emits_outgoing_and_reports_shutdown_exit_code() {
let mut s = new_session();
let first = s.on_tick(0, 1200, Some(20.0));
assert!(
!first.outgoing.is_empty(),
"the pending initial resize must be sent"
);
assert!(first.ended.is_none(), "no shutdown announced yet");
assert!(first.wait_ms <= 50, "wait is capped at 50ms");
let mut emu = ServerTerminal::new(24, 80, 0);
emu.set_exit_code(7);
let mut server = Transport::<TerminalScreen, UserInput>::new(0, 1200);
server.set_connected(true);
server.observe_rtt(20.0);
*server.current_mut() = emu.snapshot();
server.start_shutdown(0);
for dg in drive_until_nonempty(&mut server) {
s.on_datagram(10, &dg);
}
let tick = s.on_tick(10, 1200, Some(20.0));
assert_eq!(
tick.ended,
Some(Some(7)),
"a SHUTDOWN_SENTINEL remote state reports the remote shell's exit code"
);
}
#[test]
fn link_down_banner_absorbs_a_missed_keepalive_but_shows_on_a_real_stall() {
let mut s = new_session();
let mut emu = ServerTerminal::new(24, 80, 0);
emu.process(b"ready prompt $ ");
let mut server = Transport::<TerminalScreen, UserInput>::new(0, 1200);
server.set_connected(true);
server.observe_rtt(20.0);
*server.current_mut() = emu.snapshot();
for dg in drive_until_nonempty(&mut server) {
s.on_datagram(1000, &dg);
}
let absorbed = s.on_tick(1000 + 2 * crate::ssp::ACK_INTERVAL, 1200, Some(20.0));
assert!(
absorbed.status.is_none(),
"a single missed keepalive must not flash the link-down banner"
);
let boundary = s.on_tick(1000 + LINK_DOWN_GRACE_MS, 1200, Some(20.0));
assert!(
boundary.status.is_none(),
"the banner must not show until the silence exceeds the grace"
);
let stalled = s.on_tick(1000 + LINK_DOWN_GRACE_MS + 2_000, 1200, Some(20.0));
assert!(
stalled.status.is_some(),
"a silence past the grace shows the link-down banner"
);
}
#[test]
fn on_resize_resets_predictor_and_propagates() {
let mut s = new_session();
s.on_resize(40, 120);
let last_resize = s
.transport
.current()
.events()
.iter()
.rev()
.find_map(|e| match e {
InputEvent::Resize { rows, cols } => Some((*rows, *cols)),
InputEvent::Byte(_) => None,
});
assert_eq!(
last_resize,
Some((40, 120)),
"resize propagates to the server"
);
assert!(s.dirty, "a resize requires a repaint");
}
}