Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
424c2da
feat(terminal): add opt-in flow control
minpeter Aug 30, 2026
a4aeb38
Merge origin/main into feat/flow-control-631
minpeter Aug 31, 2026
96f511e
fix(terminal): harden flow-control lifecycle
minpeter Aug 31, 2026
73fb73f
fix(terminal): close flow-control review gaps
minpeter Aug 31, 2026
ae0e50d
fix(terminal): complete flow-control review
minpeter Aug 31, 2026
f586a57
fix(terminal): finish flow-control ownership review
minpeter Aug 31, 2026
b298f53
fix(terminal): close flow-control ownership gaps
minpeter Aug 31, 2026
a73be99
fix(terminal): drain retained flow-control state
minpeter Aug 31, 2026
56f3364
fix(terminal): quiesce flow-control shutdown
minpeter Aug 31, 2026
03338d2
Merge origin/main into feat/flow-control-631
minpeter Aug 31, 2026
5d868de
test(terminal): flush flow-control PTY input
minpeter Aug 31, 2026
beda35e
test(terminal): await flow-control shell readiness
minpeter Aug 31, 2026
ec57a39
test(server): synchronize jumphost response teardown
minpeter Aug 31, 2026
36e1161
fix(flow-control): close final shutdown gaps
minpeter Aug 31, 2026
c3b20da
fix(jumphost): drain buffered destination packets
minpeter Aug 31, 2026
f27049d
fix(jumphost): resume buffered output after backpressure
minpeter Aug 31, 2026
8d1e1e4
Merge origin/main into feat/flow-control-631
minpeter Aug 31, 2026
db29c03
fix(flow-control): preserve final forwarding ownership
minpeter Sep 1, 2026
57b7304
Merge origin/main into feat/flow-control-631
minpeter Sep 1, 2026
7da329f
fix(recovery): protect returning sockets through terminal HUP
minpeter Sep 1, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions .tegami/add-opt-in-flow-control.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
---
packages:
et:
type: patch
---

## Add opt-in terminal flow control

Clients can now select lossless backpressure or oldest-output discard when
terminal output outruns the network, keeping Ctrl-C and prompt responses
bounded without changing the default session behavior.
1 change: 1 addition & 0 deletions crates/et-bin/src/client_environment.rs
Original file line number Diff line number Diff line change
Expand Up @@ -189,6 +189,7 @@ mod tests {
jumphost: Some(false),
reversetunnels: vec![request; 128],
environmentvariables: Default::default(),
flowcontrol: None,
};
let mut locale = vec![
("LC_ALL".to_owned(), "C".to_owned()),
Expand Down
1 change: 1 addition & 0 deletions crates/et-bin/src/forward_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,7 @@ pub fn build(
jumphost: Some(false),
reversetunnels: reverse_tunnels,
environmentvariables: std::collections::HashMap::new(),
flowcontrol: args.flow_control.protocol_value(),
},
})
}
Expand Down
4 changes: 4 additions & 0 deletions crates/et-bin/src/terminal_protocol.rs
Original file line number Diff line number Diff line change
Expand Up @@ -182,18 +182,22 @@ mod tests {
TermInit {
environmentnames: vec!["A".to_owned()],
environmentvalues: Vec::new(),
flowcontrol: None,
},
TermInit {
environmentnames: vec!["BAD-NAME".to_owned()],
environmentvalues: vec!["value".to_owned()],
flowcontrol: None,
},
TermInit {
environmentnames: vec!["VALID".to_owned()],
environmentvalues: vec!["bad\0value".to_owned()],
flowcontrol: None,
},
TermInit {
environmentnames: vec!["VALID".to_owned()],
environmentvalues: vec!["x".repeat(MAX_ENV_VALUE + 1)],
flowcontrol: None,
},
] {
let packet = Packet::new(TerminalPacketType::TerminalInit as u8, init.encode_to_vec());
Expand Down
1 change: 1 addition & 0 deletions crates/et-bin/tests/client_bootstrap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -444,6 +444,7 @@ fn posix_client_bounds_locale_to_local_terminal_packet() {
let term_init = TermInit {
environmentnames: environment.keys().cloned().collect(),
environmentvalues: environment.values().cloned().collect(),
flowcontrol: None,
};
let packet = Packet::new(
TerminalPacketType::TerminalInit as u8,
Expand Down
148 changes: 148 additions & 0 deletions crates/et-bin/tests/flow_control_tty_qa.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,148 @@
#![cfg(unix)]
#![forbid(unsafe_code)]

mod flow_control_tty_support;

use std::fs;
use std::io::{Read, Write};
use std::sync::mpsc;
use std::thread;
use std::time::{Duration, Instant};

use flow_control_tty_support::{
receive_bytes, receive_until, Stack, ThrottleProxy, MAX_PROMPT_LATENCY, SATURATION_BYTES,
THROTTLE_BYTES_PER_SECOND,
};
use portable_pty::{native_pty_system, CommandBuilder, PtySize};

#[test]
fn flow_control_keeps_ctrl_c_and_prompt_responsive_on_a_slow_link() {
let stack = Stack::start();
let evidence = std::env::var_os("ET_FLOW_QA_EVIDENCE_DIR").map(std::path::PathBuf::from);
if let Some(directory) = &evidence {
fs::create_dir_all(directory).unwrap();
}

for mode in ["backpressure", "discard"] {
let proxy = ThrottleProxy::start(stack.port);
let pair = native_pty_system()
.openpty(PtySize {
rows: 24,
cols: 80,
pixel_width: 800,
pixel_height: 480,
})
.unwrap();
let mut client = CommandBuilder::new(env!("CARGO_BIN_EXE_et"));
client.args([
"--flow-control",
mode,
"--terminal-path",
stack.terminal.to_str().unwrap(),
"--serverfifo",
stack.router.to_str().unwrap(),
"-p",
&proxy.port.to_string(),
"127.0.0.1",
]);
client.env(
"PATH",
format!(
"{}:{}",
stack.directory.display(),
std::env::var("PATH").unwrap()
),
);
client.env("TERM", "xterm-256color");
let mut child = pair.slave.spawn_command(client).unwrap();
drop(pair.slave);

let mut writer = pair.master.take_writer().unwrap();
let mut reader = pair.master.try_clone_reader().unwrap();
let (sender, receiver) = mpsc::sync_channel(64);
let reader_thread = thread::spawn(move || {
let mut chunk = [0u8; 8192];
loop {
match reader.read(&mut chunk) {
Ok(0) | Err(_) => return,
Ok(count) if sender.send(chunk[..count].to_vec()).is_err() => return,
Ok(_) => {}
}
}
});

writer
.write_all(
b"printf 'FLOW-%s\\n' START; while :; do printf \
'0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef'; done\n",
)
.unwrap();
let output = match receive_until(
&receiver,
Vec::new(),
b"FLOW-START\r\n",
Duration::from_secs(10),
) {
Ok(output) => output,
Err(error) => {
child.kill().unwrap();
drop(writer);
let _ = child.wait();
reader_thread.join().unwrap();
panic!("{mode}: waiting for FLOW-START: {error}");
}
};
let mut output =
match receive_bytes(&receiver, output, SATURATION_BYTES, Duration::from_secs(10)) {
Ok(output) => output,
Err(error) => {
child.kill().unwrap();
drop(writer);
let _ = child.wait();
reader_thread.join().unwrap();
panic!("{mode}: saturating throttled link: {error}");
}
};
let interrupted = Instant::now();
writer
.write_all(b"\x03printf 'FLOW-%s\\n' PROMPT\n")
.unwrap();
output = match receive_until(&receiver, output, b"FLOW-PROMPT\r\n", MAX_PROMPT_LATENCY) {
Ok(output) => output,
Err(error) => {
child.kill().unwrap();
drop(writer);
let _ = child.wait();
reader_thread.join().unwrap();
panic!("{mode}: waiting for Ctrl-C prompt: {error}");
}
};
let latency = interrupted.elapsed();
assert!(
latency <= MAX_PROMPT_LATENCY,
"{mode} Ctrl-C-to-prompt latency {latency:?} exceeded {MAX_PROMPT_LATENCY:?}"
);

child.kill().unwrap();
drop(writer);
while let Ok(chunk) = receiver.recv_timeout(Duration::from_millis(100)) {
output.extend(chunk);
}
let _ = child.wait();
reader_thread.join().unwrap();

if let Some(directory) = &evidence {
fs::write(directory.join(format!("{mode}.ansi")), &output).unwrap();
fs::write(
directory.join(format!("{mode}.json")),
format!(
"{{\"mode\":\"{mode}\",\"rate_bytes_per_second\":{THROTTLE_BYTES_PER_SECOND},\
\"saturation_bytes\":{SATURATION_BYTES},\"ctrl_c_prompt_millis\":{},\
\"pass\":true}}\n",
latency.as_millis()
),
)
.unwrap();
}
}
}
178 changes: 178 additions & 0 deletions crates/et-bin/tests/flow_control_tty_support.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,178 @@
#![cfg(unix)]
#![allow(dead_code)]

use std::fs;
use std::io::{self, BufRead, BufReader, Read, Write};
use std::net::{Ipv4Addr, Shutdown, TcpListener, TcpStream};
use std::os::unix::fs::{symlink, PermissionsExt};
use std::process::{Command, Stdio};
use std::sync::mpsc;
use std::thread;
use std::time::{Duration, Instant};

use nix::sys::signal::{kill, Signal};
use nix::unistd::Pid;
use wait_timeout::ChildExt;

const TIMEOUT: Duration = Duration::from_secs(10);
pub const MAX_PROMPT_LATENCY: Duration = Duration::from_secs(5);
pub const THROTTLE_BYTES_PER_SECOND: usize = 100 * 1024;
pub const SATURATION_BYTES: usize = 128 * 1024;

pub fn receive_until(
receiver: &mpsc::Receiver<Vec<u8>>,
mut output: Vec<u8>,
marker: &[u8],
timeout: Duration,
) -> Result<Vec<u8>, mpsc::RecvTimeoutError> {
let deadline = Instant::now() + timeout;
while !output.windows(marker.len()).any(|window| window == marker) {
let Some(remaining) = deadline.checked_duration_since(Instant::now()) else {
return Err(mpsc::RecvTimeoutError::Timeout);
};
output.extend(receiver.recv_timeout(remaining)?);
}
Ok(output)
}

pub fn receive_bytes(
receiver: &mpsc::Receiver<Vec<u8>>,
mut output: Vec<u8>,
additional: usize,
timeout: Duration,
) -> Result<Vec<u8>, mpsc::RecvTimeoutError> {
let target = output.len() + additional;
let deadline = Instant::now() + timeout;
while output.len() < target {
let Some(remaining) = deadline.checked_duration_since(Instant::now()) else {
return Err(mpsc::RecvTimeoutError::Timeout);
};
output.extend(receiver.recv_timeout(remaining)?);
}
Ok(output)
}

pub struct ThrottleProxy {
pub port: u16,
}

impl ThrottleProxy {
pub fn start(server_port: u16) -> Self {
let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).unwrap();
let port = listener.local_addr().unwrap().port();
let _worker: thread::JoinHandle<io::Result<()>> = thread::spawn(move || {
let (mut client, _) = listener.accept()?;
let mut server = TcpStream::connect((Ipv4Addr::LOCALHOST, server_port))?;
let mut client_read = client.try_clone()?;
let mut server_write = server.try_clone()?;
let upstream = thread::spawn(move || io::copy(&mut client_read, &mut server_write));

let started = Instant::now();
let mut transferred = 0usize;
let mut chunk = [0u8; 8192];
loop {
let count = server.read(&mut chunk)?;
if count == 0 {
break;
}
client.write_all(&chunk[..count])?;
transferred += count;
let expected =
Duration::from_secs_f64(transferred as f64 / THROTTLE_BYTES_PER_SECOND as f64);
if let Some(remaining) = expected.checked_sub(started.elapsed()) {
thread::sleep(remaining);
}
}
let _ = client.shutdown(Shutdown::Both);
let _ = server.shutdown(Shutdown::Both);
upstream
.join()
.map_err(|_| io::Error::other("proxy upload thread panicked"))??;
Ok(())
});
Self { port }
}
}

pub struct Stack {
pub directory: std::path::PathBuf,
pub router: std::path::PathBuf,
pub terminal: std::path::PathBuf,
pub port: u16,
server: std::process::Child,
}

impl Stack {
pub fn start() -> Self {
let directory =
std::env::temp_dir().join(format!("et-rs-flow-control-qa-{}", std::process::id()));
let _ = fs::remove_dir_all(&directory);
fs::create_dir(&directory).unwrap();
fs::set_permissions(&directory, fs::Permissions::from_mode(0o700)).unwrap();
let router = directory.join("router.sock");
let reserved = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).unwrap();
let port = reserved.local_addr().unwrap().port();
drop(reserved);
let config = directory.join("et.cfg");
fs::write(
&config,
format!(
"[Networking]\nport={port}\nbind_ip=127.0.0.1\n[Debug]\nserverfifo={}\n",
router.display()
),
)
.unwrap();
let mut server = Command::new(env!("CARGO_BIN_EXE_et"))
.args(["server", "--cfgfile"])
.arg(&config)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap();
wait_ready(&mut server, port, &router);
let ssh = directory.join("ssh");
fs::write(
&ssh,
"#!/bin/sh\nif [ \"$1\" = \"-G\" ]; then\n\
printf 'hostname 127.0.0.1\\nuser tester\\n'; exit 0; fi\n\
for last do :; done\nexec /bin/sh -c \"$last\"\n",
)
.unwrap();
fs::set_permissions(&ssh, fs::Permissions::from_mode(0o755)).unwrap();
let terminal = directory.join("etterminal");
symlink(env!("CARGO_BIN_EXE_et"), &terminal).unwrap();
Self {
directory,
router,
terminal,
port,
server,
}
}
}

impl Drop for Stack {
fn drop(&mut self) {
let pid = Pid::from_raw(i32::try_from(self.server.id()).unwrap());
let _ = kill(pid, Signal::SIGTERM);
let _ = self.server.wait_timeout(TIMEOUT);
let _ = fs::remove_dir_all(&self.directory);
}
}

fn wait_ready(server: &mut std::process::Child, port: u16, router: &std::path::Path) {
let stdout = server.stdout.take().unwrap();
let (sender, receiver) = mpsc::sync_channel(1);
thread::spawn(move || {
let mut line = String::new();
let result = BufReader::new(stdout).read_line(&mut line).map(|_| line);
let _ = sender.send(result);
});
assert_eq!(
receiver.recv_timeout(TIMEOUT).unwrap().unwrap(),
format!(
"ETSERVER_READY tcp=127.0.0.1:{port} router={}\n",
router.display()
)
);
}
Loading
Loading