mod audit;
pub mod cli;
pub mod session;
pub use cli::{serve, ServeArgs};
use std::time::Duration;
use crate::input::{UserInput, WireEvent};
use crate::ssp::{RecvOutcome, Transport};
use crate::terminal::TerminalScreen;
use crate::transport_iroh::{IrohChannel, MonoClock};
use session::SharedSession;
use tracing::{info, warn};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SessionExit {
Detached,
ShellExited,
}
#[derive(Default, Clone, Copy, PartialEq, Eq)]
enum Ss3State {
#[default]
Ground,
Esc,
Ss3,
}
#[derive(Default)]
struct CursorKeyNormalizer {
state: Ss3State,
}
impl CursorKeyNormalizer {
fn normalize(&mut self, input: &[u8], app_cursor: bool) -> Vec<u8> {
let mut out = Vec::with_capacity(input.len() + 1);
for &b in input {
match self.state {
Ss3State::Ground => {
if b == 0x1b {
self.state = Ss3State::Esc;
}
out.push(b); }
Ss3State::Esc => {
if b == b'O' {
self.state = Ss3State::Ss3; } else {
self.state = Ss3State::Ground;
out.push(b);
}
}
Ss3State::Ss3 => {
self.state = Ss3State::Ground;
out.push(if !app_cursor && (b'A'..=b'D').contains(&b) {
b'['
} else {
b'O'
});
out.push(b);
}
}
}
out
}
}
fn coalesce_drained_input(
input_diff: &[WireEvent],
cursor_keys: &mut CursorKeyNormalizer,
app_cursor: bool,
) -> (Vec<u8>, Option<(u16, u16)>) {
let mut keys = Vec::new();
let mut last_resize: Option<(u16, u16)> = None;
for w in input_diff {
match w {
WireEvent::Keys(b) => keys.extend(cursor_keys.normalize(b, app_cursor)),
WireEvent::Resize { rows, cols } => {
last_resize = Some(crate::terminal::clamp_dims(*rows, *cols));
}
}
}
(keys, last_resize)
}
struct DrainedInput {
keys: Vec<u8>,
resize: Option<(u16, u16)>,
frame: u64,
}
struct ServerSession {
transport: Transport<TerminalScreen, UserInput>,
cursor_keys: CursorKeyNormalizer,
dirty: bool,
}
impl ServerSession {
fn new(now: u64, mtu: usize) -> Self {
let mut transport = Transport::<TerminalScreen, UserInput>::new(now, mtu);
transport.set_connected(true);
Self {
transport,
cursor_keys: CursorKeyNormalizer::default(),
dirty: true, }
}
fn observe_link(&mut self, mtu: usize, rtt_ms: Option<f64>) {
self.transport.set_mtu(mtu);
if let Some(rtt) = rtt_ms {
self.transport.observe_rtt(rtt);
}
}
const fn needs_snapshot(&self, echo_changed: bool) -> bool {
self.dirty || echo_changed
}
fn install_snapshot(&mut self, snapshot: Option<TerminalScreen>) {
if let Some(screen) = snapshot {
*self.transport.current_mut() = screen;
}
self.dirty = false;
}
fn wait_ms(&mut self, now: u64, echo_wait: u64) -> u64 {
self.transport.wait_time(now).min(echo_wait).min(1000)
}
fn mark_dirty(&mut self) {
self.dirty = true;
}
fn recv(&mut self, now: u64, bytes: &[u8]) -> RecvOutcome {
self.transport.recv(now, bytes)
}
fn drain_input(&mut self, app_cursor: bool) -> Option<DrainedInput> {
let diff = self.transport.get_remote_diff();
if diff.is_empty() {
return None;
}
let frame = self.transport.remote_num();
let (keys, resize) = coalesce_drained_input(&diff, &mut self.cursor_keys, app_cursor);
Some(DrainedInput {
keys,
resize,
frame,
})
}
fn tick(&mut self, now: u64, child_alive: bool) -> Vec<Vec<u8>> {
if !child_alive && !self.transport.shutdown_in_progress() {
self.transport.start_shutdown(now);
}
self.transport.tick(now)
}
fn shutdown_complete(&self, now: u64) -> bool {
self.transport.shutdown_in_progress()
&& (self.transport.shutdown_acknowledged()
|| self.transport.shutdown_ack_timed_out(now))
}
}
pub async fn run_attached(
conn: iroh::endpoint::Connection,
handle: SharedSession,
) -> anyhow::Result<SessionExit> {
let channel = IrohChannel::new(conn);
let clock = MonoClock::new();
let mut session = ServerSession::new(clock.now_ms(), channel.max_datagram_size());
loop {
let now = clock.now_ms();
session.observe_link(channel.max_datagram_size(), channel.rtt_ms());
let (echo_wait, child_alive) = {
let mut s = handle.session.lock().await;
let echo_changed = s.emu.set_echo_ack(now);
let snapshot = session
.needs_snapshot(echo_changed)
.then(|| s.emu.snapshot());
session.install_snapshot(snapshot);
(s.emu.echo_ack_wait_time(now), s.child_alive)
};
let sleep_ms = session.wait_ms(now, echo_wait);
tokio::select! {
_ = handle.changed.notified() => session.mark_dirty(),
dg = channel.recv() => {
match dg {
Ok(bytes) => {
let now = clock.now_ms();
if session.recv(now, &bytes) == RecvOutcome::NewState {
let mut s = handle.session.lock().await;
let app_cursor = s.emu.application_cursor();
if let Some(input) = session.drain_input(app_cursor) {
if !input.keys.is_empty() {
if let Err(e) = s.pty.write_input(&input.keys) {
warn!(error = %e, "pty write failed");
}
}
let resized = input.resize.is_some();
if let Some((rows, cols)) = input.resize {
if let Err(e) = s.pty.resize(rows, cols) {
warn!(error = %e, rows, cols, "pty resize failed");
}
s.emu.resize(rows, cols);
}
s.emu.register_input_frame(input.frame, now);
drop(s);
if resized {
session.mark_dirty();
}
}
}
}
Err(e) => {
info!(reason = %e, "connection closed by peer (detaching)");
channel.close(0, b"client detached");
return Ok(SessionExit::Detached);
}
}
}
_ = tokio::time::sleep(Duration::from_millis(sleep_ms)) => {}
}
let now = clock.now_ms();
for datagram in session.tick(now, child_alive) {
channel.send(&datagram);
}
if session.shutdown_complete(now) {
channel.close(0, b"session ended");
return Ok(SessionExit::ShellExited);
}
}
}
pub async fn run_session(
conn: iroh::endpoint::Connection,
shell: Option<String>,
scrollback: usize,
) -> anyhow::Result<()> {
let handle = session::spawn_session(shell.as_deref(), scrollback)?;
let _ = run_attached(conn, handle.clone()).await?;
let _ = handle.session.lock().await.pty.kill();
Ok(())
}
#[cfg(test)]
mod tests {
use super::{coalesce_drained_input, CursorKeyNormalizer, ServerSession};
use crate::input::{UserInput, WireEvent};
use crate::ssp::{RecvOutcome, Transport};
use crate::terminal::TerminalScreen;
fn norm(chunks: &[&[u8]], app_cursor: bool) -> Vec<u8> {
let mut n = CursorKeyNormalizer::default();
let mut out = Vec::new();
for c in chunks {
out.extend(n.normalize(c, app_cursor));
}
out
}
#[test]
fn ss3_arrows_rewrite_to_csi_when_not_in_application_cursor_mode() {
assert_eq!(norm(&[b"\x1bOA"], false), b"\x1b[A");
assert_eq!(norm(&[b"\x1bOD"], false), b"\x1b[D");
}
#[test]
fn ss3_arrows_preserved_in_application_cursor_mode() {
assert_eq!(norm(&[b"\x1bOA"], true), b"\x1bOA");
}
#[test]
fn csi_arrows_and_plain_bytes_pass_through() {
assert_eq!(norm(&[b"\x1b[A"], false), b"\x1b[A");
assert_eq!(norm(&[b"ls\r"], false), b"ls\r");
assert_eq!(norm(&[b"\x1bi"], false), b"\x1bi");
}
#[test]
fn ss3_sequence_split_across_chunks_normalizes() {
assert_eq!(norm(&[b"\x1b", b"O", b"A"], false), b"\x1b[A");
assert_eq!(norm(&[b"\x1b", b"[", b"A"], false), b"\x1b[A");
}
#[test]
fn coalesce_keeps_only_the_last_resize_and_concatenates_keys() {
let mut norm = CursorKeyNormalizer::default();
let diff = vec![
WireEvent::Keys(b"ab".to_vec()),
WireEvent::Resize { rows: 10, cols: 20 },
WireEvent::Keys(b"cd".to_vec()),
WireEvent::Resize { rows: 30, cols: 40 },
WireEvent::Resize {
rows: 65000,
cols: 1,
}, WireEvent::Keys(b"ef".to_vec()),
];
let (keys, last_resize) = coalesce_drained_input(&diff, &mut norm, false);
assert_eq!(keys, b"abcdef", "keystrokes concatenate in order");
assert_eq!(
last_resize,
Some(crate::terminal::clamp_dims(65000, 1)),
"only the final resize survives, clamped to [MIN_DIM, MAX_DIM]"
);
}
#[test]
fn coalesce_with_no_resize_returns_none() {
let mut norm = CursorKeyNormalizer::default();
let diff = vec![WireEvent::Keys(b"x".to_vec())];
let (keys, last_resize) = coalesce_drained_input(&diff, &mut norm, false);
assert_eq!(keys, b"x");
assert!(last_resize.is_none(), "no resize event -> None");
}
#[test]
fn server_session_snapshot_gating() {
let mut s = ServerSession::new(0, 1200);
assert!(s.needs_snapshot(false), "the first pass always snapshots");
s.install_snapshot(Some(TerminalScreen::default()));
assert!(!s.needs_snapshot(false), "clean after a snapshot");
assert!(
s.needs_snapshot(true),
"an echo-ack advance forces a snapshot even when clean (else confirmations stall)"
);
s.mark_dirty();
assert!(
s.needs_snapshot(false),
"a changed-pulse / applied resize re-arms the snapshot"
);
}
#[test]
fn server_session_shutdown_handshake_progresses() {
let mut s = ServerSession::new(0, 1200);
let _ = s.tick(0, true); assert!(!s.shutdown_complete(0));
let _ = s.tick(10, false); assert!(
!s.shutdown_complete(10),
"shutdown just started: neither acked nor timed out yet"
);
assert!(
s.shutdown_complete(10_000_000),
"far in the future the unacked shutdown times out -> reapable"
);
}
#[test]
fn server_session_drains_coalesced_input_from_a_real_datagram() {
let mut client = Transport::<UserInput, TerminalScreen>::new(0, 1200);
client.set_connected(true);
client.current_mut().push_bytes(b"ls\r");
client.current_mut().push_resize(10, 20);
client.current_mut().push_resize(30, 40);
let datagrams = client.tick(1000);
assert!(
!datagrams.is_empty(),
"the client transmits its queued input"
);
let mut server = ServerSession::new(0, 1200);
let mut drained = None;
for dg in &datagrams {
if server.recv(1000, dg) == RecvOutcome::NewState {
drained = server.drain_input(false);
}
}
let input = drained.expect("the server drained the client's input");
assert_eq!(
input.keys, b"ls\r",
"keystrokes concatenate in order through the normalizer"
);
assert_eq!(
input.resize,
Some((30, 40)),
"KOH-05: only the final resize survives (clamped)"
);
}
}