use crate::{Continuation, Controlled, ControlledChild, Controller, ControllerReader,
ControllerWriter, LINE_FEED_BYTE, PtyCommandBuilder, PtyInputEvent,
PtyReadWriteOutputEvent, PtyReadWriteSession, READ_BUFFER_SIZE, ok,
pty_common_io::{create_pty_pair, spawn_command_in_pty}};
use miette::{IntoDiagnostic, miette};
use portable_pty::PtySize;
use std::{io::{Read, Write},
sync::mpsc::RecvTimeoutError,
time::Duration};
impl PtyCommandBuilder {
pub fn spawn_read_write(
self,
pty_size: PtySize,
) -> miette::Result<PtyReadWriteSession> {
let (
output_evt_ch_tx_half,
output_evt_ch_rx_half,
) = tokio::sync::mpsc::unbounded_channel::<PtyReadWriteOutputEvent>();
let (
input_evt_ch_tx_half,
input_evt_ch_rx_half,
) = tokio::sync::mpsc::unbounded_channel::<PtyInputEvent>();
let command = self.build()?;
let (controller, controlled): (Controller, Controlled) =
create_pty_pair(pty_size)?;
let mut controlled_child: ControlledChild =
spawn_command_in_pty(&controlled, command)?;
let child_process_terminate_handle = controlled_child.clone_killer();
let session_completion_handle = tokio::spawn(async move {
let output_reader_task_handle = {
let controller_reader = controller
.try_clone_reader()
.map_err(|e| miette!("Failed to clone reader: {}", e))?;
spawn_blocking_passthrough_with_mode_detection_reader_task(
controller_reader,
output_evt_ch_tx_half.clone(),
)
};
let input_writer_task_handle = create_controller_input_writer_task(
controller,
input_evt_ch_rx_half,
output_evt_ch_tx_half.clone(),
);
let status = tokio::task::spawn_blocking(move || controlled_child.wait())
.await
.into_diagnostic()?
.into_diagnostic()?;
let exit_code = status.exit_code();
let _unused =
output_evt_ch_tx_half.send(PtyReadWriteOutputEvent::Exit(status));
drop(controlled);
let _unused = input_writer_task_handle.await;
let _unused = output_reader_task_handle.await;
Ok(portable_pty::ExitStatus::with_exit_code(exit_code))
});
Ok(PtyReadWriteSession {
input_event_ch_tx_half: input_evt_ch_tx_half,
output_event_receiver_half: output_evt_ch_rx_half,
pinned_boxed_session_completion_handle: Box::pin(session_completion_handle),
child_process_terminate_handle,
})
}
}
#[must_use]
fn create_controller_input_writer_task(
controller: Controller,
input_evt_ch_rx_half: tokio::sync::mpsc::UnboundedReceiver<PtyInputEvent>,
output_evt_ch_tx_half: tokio::sync::mpsc::UnboundedSender<PtyReadWriteOutputEvent>,
) -> tokio::task::JoinHandle<miette::Result<()>> {
let (
input_evt_bridge_sync_tx_half,
input_evt_bridge_sync_rx_half,
) = std::sync::mpsc::channel::<PtyInputEvent>();
let input_writer_task_handle = spawn_blocking_writer_task(
controller,
input_evt_bridge_sync_rx_half,
output_evt_ch_tx_half.clone(),
);
let input_writer_bridge_handle = spawn_async_to_sync_bridge_task(
input_evt_ch_rx_half,
input_evt_bridge_sync_tx_half,
);
tokio::spawn(async move {
let (_bridge, writer) =
tokio::join!(input_writer_bridge_handle, input_writer_task_handle);
writer.map_err(|e| miette!("Input writer task failed: {}", e))?
})
}
#[must_use]
fn spawn_async_to_sync_bridge_task(
mut input_evt_ch_rx_half: tokio::sync::mpsc::UnboundedReceiver<PtyInputEvent>,
input_evt_bridge_sync_tx_half: std::sync::mpsc::Sender<PtyInputEvent>,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
while let Some(input) = input_evt_ch_rx_half.recv().await {
if input_evt_bridge_sync_tx_half.send(input).is_err() {
break;
}
}
let _unused = input_evt_bridge_sync_tx_half.send(PtyInputEvent::Close);
})
}
#[must_use]
fn spawn_blocking_writer_task(
controller: Controller,
input_evt_bridge_sync_rx_half: std::sync::mpsc::Receiver<PtyInputEvent>,
output_evt_ch_tx_half: tokio::sync::mpsc::UnboundedSender<PtyReadWriteOutputEvent>,
) -> tokio::task::JoinHandle<miette::Result<()>> {
tokio::task::spawn_blocking(move || -> miette::Result<()> {
let mut writer = controller
.take_writer()
.map_err(|e| miette!("Failed to take PTY writer: {}", e))?;
loop {
match input_evt_bridge_sync_rx_half.recv_timeout(Duration::from_millis(100)) {
Err(RecvTimeoutError::Disconnected) => {
break;
}
Err(RecvTimeoutError::Timeout) => {
}
Ok(input) => {
match handle_pty_input_event(
input,
&mut writer,
&controller,
&output_evt_ch_tx_half,
)? {
Continuation::Continue => {
}
Continuation::Stop => {
break;
}
}
}
}
}
drop(controller);
Ok(())
})
}
fn handle_pty_input_event(
input: PtyInputEvent,
writer: &mut ControllerWriter,
controller: &Controller,
output_evt_ch_tx_half: &tokio::sync::mpsc::UnboundedSender<PtyReadWriteOutputEvent>,
) -> miette::Result<Continuation> {
match input {
PtyInputEvent::Write(bytes) => write_to_pty_with_flush(
writer,
&bytes,
"Failed to write to PTY",
output_evt_ch_tx_half,
)?,
PtyInputEvent::WriteLine(text) => write_to_pty_with_flush(
writer,
&{
let mut data = text.into_bytes();
data.push(LINE_FEED_BYTE);
data
},
"Failed to write line to PTY",
output_evt_ch_tx_half,
)?,
PtyInputEvent::SendControl(ctrl, mode) => write_to_pty_with_flush(
writer,
&ctrl.to_bytes(mode),
"Failed to send control char to PTY",
output_evt_ch_tx_half,
)?,
PtyInputEvent::Resize(size) => controller.resize(size).map_err(|e| {
let _unused = output_evt_ch_tx_half.send(
PtyReadWriteOutputEvent::WriteError(miette!("Resize failed: {e}")),
);
miette!("Failed to resize PTY")
})?,
PtyInputEvent::Flush => writer.flush().map_err(|e| {
let _unused = output_evt_ch_tx_half.send(
PtyReadWriteOutputEvent::WriteError(miette!("Flush failed: {e}")),
);
miette!("Failed to flush PTY")
})?,
PtyInputEvent::Close => return Ok(Continuation::Stop),
}
Ok(Continuation::Continue)
}
fn write_to_pty_with_flush(
writer: &mut ControllerWriter,
data: &[u8],
error_msg: &str,
output_evt_ch_tx_half: &tokio::sync::mpsc::UnboundedSender<PtyReadWriteOutputEvent>,
) -> miette::Result<()> {
writer.write_all(data).map_err(|e| {
let _unused = output_evt_ch_tx_half.send(PtyReadWriteOutputEvent::WriteError(
miette!("Write failed: {}", e),
));
miette!("{error_msg}")
})?;
writer.flush().map_err(|e| {
let _unused = output_evt_ch_tx_half.send(PtyReadWriteOutputEvent::WriteError(
miette!("Flush failed: {}", e),
));
miette!("{error_msg}")
})?;
ok!()
}
fn spawn_blocking_passthrough_with_mode_detection_reader_task(
mut controller_reader: ControllerReader,
output_evt_ch_tx_half: tokio::sync::mpsc::UnboundedSender<PtyReadWriteOutputEvent>,
) -> tokio::task::JoinHandle<miette::Result<()>> {
tokio::task::spawn_blocking(move || -> miette::Result<()> {
let mut read_buffer = [0u8; READ_BUFFER_SIZE];
let mut mode_detector =
crate::pty_core::pty_output_events::CursorModeDetector::new();
loop {
match controller_reader.read(&mut read_buffer) {
Ok(0) => {
let _unused = output_evt_ch_tx_half.send(
PtyReadWriteOutputEvent::UnexpectedExit(
"PTY closed (EOF)".to_string(),
),
);
break;
}
Err(e) => {
let _unused = output_evt_ch_tx_half.send(
PtyReadWriteOutputEvent::UnexpectedExit(format!(
"Read error: {e}"
)),
);
break;
}
Ok(n) => {
let data = &read_buffer[..n];
if let Some(new_mode) = mode_detector.scan_for_mode_change(data) {
let _unused = output_evt_ch_tx_half
.send(PtyReadWriteOutputEvent::CursorModeChange(new_mode));
}
let _unused = output_evt_ch_tx_half
.send(PtyReadWriteOutputEvent::Output(data.to_vec()));
}
}
}
drop(controller_reader);
Ok(())
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{ControlSequence, CursorKeyMode, try_create_temp_dir};
use tokio::sync::mpsc::unbounded_channel;
#[test]
fn test_all_pty_read_write_in_isolated_process() {
if is_ci::uncached() {
println!(
"Skipping PTY tests in CI environment due to PTY resource limitations"
);
return;
}
if let Ok(test_name) = std::env::var("ISOLATED_PTY_SINGLE_TEST") {
run_single_pty_test_by_name(&test_name);
std::process::exit(0);
}
let tests = vec![
"test_simple_command_lifecycle",
"test_cat_with_input",
#[cfg(not(target_os = "windows"))]
"test_shell_calculation",
#[cfg(not(target_os = "windows"))]
"test_shell_echo_output",
"test_multiple_control_characters",
"test_raw_escape_sequences",
#[cfg(not(target_os = "windows"))]
"test_htop_interactive_with_cursor_modes",
];
let mut failed_tests = Vec::new();
for &test_name in &tests {
println!("Running {test_name} in isolated process...");
let output = run_single_pty_test_in_isolated_process(test_name);
let stderr = String::from_utf8_lossy(&output.stderr).to_string();
let stdout = String::from_utf8_lossy(&output.stdout).to_string();
if !output.status.success()
|| stderr.contains("panicked at")
|| stderr.contains("Test failed with error")
{
failed_tests.push(test_name);
eprintln!("❌ {test_name} failed:");
eprintln!(" Exit status: {:?}", output.status);
eprintln!(" Stdout: {stdout}");
eprintln!(" Stderr: {stderr}");
} else {
println!("✅ {test_name} passed");
}
}
if !failed_tests.is_empty() {
eprintln!("⚠️ The following PTY tests failed: {failed_tests:?}");
eprintln!(
"This may be due to PTY environment limitations in the test environment."
);
eprintln!(
"PTY tests can be sensitive to system resources, configuration, and CI environments."
);
if failed_tests.len() > tests.len() / 2 {
panic!(
"Too many PTY tests failed ({}/{}). This indicates a serious PTY system issue.",
failed_tests.len(),
tests.len()
);
} else {
println!(
"Continuing despite {} PTY test failures - this is acceptable for environment-sensitive tests.",
failed_tests.len()
);
}
}
println!("All PTY read-write tests completed successfully in isolated processes");
}
fn run_single_pty_test_in_isolated_process(test_name: &str) -> std::process::Output {
let temp_dir =
try_create_temp_dir().expect("Failed to create temp dir for PTY test");
let current_exe = std::env::current_exe().unwrap();
let mut cmd = std::process::Command::new(¤t_exe);
cmd.current_dir(temp_dir.as_ref())
.env("ISOLATED_PTY_SINGLE_TEST", test_name)
.env("RUST_BACKTRACE", "1")
.args([
"--test-threads",
"1",
"test_all_pty_read_write_in_isolated_process",
]);
cmd.output().expect("Failed to run isolated PTY test")
}
#[allow(clippy::missing_errors_doc)]
fn run_single_pty_test_by_name(test_name: &str) {
let runtime = tokio::runtime::Runtime::new()
.expect("Failed to create Tokio runtime for PTY test");
runtime.block_on(async {
let result = match test_name {
"test_simple_command_lifecycle" => test_simple_command_lifecycle().await,
"test_cat_with_input" => test_cat_with_input().await,
#[cfg(not(target_os = "windows"))]
"test_shell_calculation" => test_shell_calculation().await,
#[cfg(not(target_os = "windows"))]
"test_shell_echo_output" => test_shell_echo_output().await,
"test_multiple_control_characters" => {
test_multiple_control_characters().await
}
"test_raw_escape_sequences" => test_raw_escape_sequences().await,
#[cfg(not(target_os = "windows"))]
"test_htop_interactive_with_cursor_modes" => {
test_htop_interactive_with_cursor_modes().await
}
_ => panic!("Unknown test name: {test_name}"),
};
if let Err(e) = result {
panic!("{test_name} failed: {e}");
}
println!("{test_name} completed successfully!");
});
}
async fn test_simple_command_lifecycle() -> miette::Result<()> {
use tokio::time::timeout;
let temp_dir = std::env::temp_dir();
let mut session = PtyCommandBuilder::new("echo")
.args(["Hello, PTY!"])
.cwd(temp_dir)
.spawn_read_write(PtySize::default())
.unwrap();
tokio::time::sleep(Duration::from_millis(50)).await;
let mut output = String::new();
let mut events_received = Vec::new();
let mut saw_exit = false;
let result = timeout(Duration::from_secs(10), async {
while let Some(event) = session.output_event_receiver_half.recv().await {
match &event {
PtyReadWriteOutputEvent::Output(data) => {
let data_str = String::from_utf8_lossy(data);
output.push_str(&data_str);
events_received.push(format!(
"Output({} bytes): '{}'",
data.len(),
data_str
));
}
PtyReadWriteOutputEvent::Exit(status) => {
saw_exit = true;
events_received.push(format!("Exit({status:?})"));
assert!(
status.success(),
"Command should succeed with status: {status:?}"
);
break;
}
other => {
events_received.push(format!("{other:?}"));
}
}
}
})
.await;
assert!(
result.is_ok(),
"Test timed out after 10 seconds. Events received: {events_received:?}, Output so far: '{output}'"
);
assert!(
saw_exit,
"Should see exit event. Events received: {events_received:?}, Output: '{output}'"
);
assert!(
output.contains("Hello, PTY!"),
"Output should contain 'Hello, PTY!'. Events received: {events_received:?}, Full output was: '{output}'"
);
Ok(())
}
async fn test_cat_with_input() -> miette::Result<()> {
use tokio::time::timeout;
let temp_dir = std::env::temp_dir();
let mut session = PtyCommandBuilder::new("cat")
.cwd(temp_dir)
.spawn_read_write(PtySize::default())
.unwrap();
session
.input_event_ch_tx_half
.send(PtyInputEvent::WriteLine("test input".into()))
.unwrap();
session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::CtrlD,
CursorKeyMode::default(),
))
.unwrap();
let mut output = String::new();
let mut events_received = Vec::new();
let mut saw_exit = false;
let result = timeout(Duration::from_secs(10), async {
while let Some(event) = session.output_event_receiver_half.recv().await {
match &event {
PtyReadWriteOutputEvent::Output(data) => {
let data_str = String::from_utf8_lossy(data);
output.push_str(&data_str);
events_received.push(format!(
"Output({} bytes): '{}'",
data.len(),
data_str
));
}
PtyReadWriteOutputEvent::Exit(status) => {
saw_exit = true;
events_received.push(format!("Exit({status:?})"));
assert!(
status.success(),
"Cat should succeed with status: {status:?}"
);
break;
}
other => {
events_received.push(format!("{other:?}"));
}
}
}
})
.await;
assert!(
result.is_ok(),
"Test timed out after 10 seconds. Events received: {events_received:?}, Output so far: '{output}'"
);
assert!(
saw_exit,
"Should see exit event. Events received: {events_received:?}, Output: '{output}'"
);
assert!(
output.contains("test input"),
"Output should contain 'test input'. Events received: {events_received:?}, Full output was: '{output}'"
);
Ok(())
}
#[cfg(not(target_os = "windows"))]
async fn test_shell_calculation() -> miette::Result<()> {
use tokio::time::timeout;
let temp_dir = std::env::temp_dir();
let mut session = PtyCommandBuilder::new("sh")
.args(["-c", "echo $((2+3)); echo 'Hello from Shell'"])
.cwd(temp_dir)
.spawn_read_write(PtySize::default())
.unwrap();
let mut output = String::new();
let mut events_received = Vec::new();
let mut saw_exit = false;
let result = timeout(Duration::from_secs(10), async {
while let Some(event) = session.output_event_receiver_half.recv().await {
match &event {
PtyReadWriteOutputEvent::Output(data) => {
let data_str = String::from_utf8_lossy(data);
output.push_str(&data_str);
events_received.push(format!(
"Output({} bytes): '{}'",
data.len(),
data_str
));
}
PtyReadWriteOutputEvent::Exit(status) => {
saw_exit = true;
events_received.push(format!("Exit({status:?})"));
break;
}
other => {
events_received.push(format!("{other:?}"));
}
}
}
})
.await;
assert!(
result.is_ok(),
"Shell session timed out after 10 seconds. Events received: {events_received:?}, Output so far: '{output}'"
);
assert!(
saw_exit,
"Should see exit event. Events received: {events_received:?}, Output: '{output}'"
);
assert!(
output.contains('5'),
"Should see result of 2+3. Events received: {events_received:?}, Full output was: '{output}'"
);
assert!(
output.contains("Hello from Shell"),
"Should see hello message. Events received: {events_received:?}, Full output was: '{output}'"
);
Ok(())
}
#[cfg(not(target_os = "windows"))]
async fn test_shell_echo_output() -> miette::Result<()> {
use tokio::time::timeout;
let mut session = PtyCommandBuilder::new("/bin/echo")
.args(["Test output from echo"])
.spawn_read_write(PtySize::default())
.unwrap();
tokio::time::sleep(Duration::from_millis(100)).await;
let mut output = String::new();
let mut events_received = Vec::new();
let mut saw_exit = false;
let result = timeout(Duration::from_secs(10), async {
while let Some(event) = session.output_event_receiver_half.recv().await {
match &event {
PtyReadWriteOutputEvent::Output(data) => {
let data_str = String::from_utf8_lossy(data);
output.push_str(&data_str);
events_received.push(format!(
"Output({} bytes): '{}'",
data.len(),
data_str
));
}
PtyReadWriteOutputEvent::Exit(status) => {
saw_exit = true;
events_received.push(format!("Exit({status:?})"));
assert!(
status.success(),
"Shell should succeed with status: {status:?}"
);
break;
}
other => {
events_received.push(format!("{other:?}"));
}
}
}
})
.await;
assert!(
result.is_ok(),
"Test timed out after 10 seconds. Events received: {events_received:?}, Output so far: '{output}'"
);
assert!(
saw_exit,
"Should see exit event. Events received: {events_received:?}, Output: '{output}'"
);
assert!(
output.contains("Test output from echo"),
"Should see echo output. Events received: {events_received:?}, Full output was: '{output}'"
);
Ok(())
}
async fn test_multiple_control_characters() -> miette::Result<()> {
use tokio::time::timeout;
let mut session = PtyCommandBuilder::new("cat")
.spawn_read_write(PtySize::default())
.unwrap();
session
.input_event_ch_tx_half
.send(PtyInputEvent::WriteLine("Test line".into()))
.unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
if session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::Enter,
CursorKeyMode::default(),
))
.is_err()
{
let mut output = String::new();
while let Some(event) = session.output_event_receiver_half.recv().await {
match event {
PtyReadWriteOutputEvent::Output(data) => {
output.push_str(&String::from_utf8_lossy(&data));
}
PtyReadWriteOutputEvent::Exit(status) => {
panic!(
"PTY exited early with status: {status:?}, output: '{output}'"
);
}
_ => {}
}
}
panic!("Failed to send Enter control character, output so far: '{output}'");
}
tokio::time::sleep(Duration::from_millis(10)).await;
session
.input_event_ch_tx_half
.send(PtyInputEvent::Write(b"No newline".to_vec()))
.unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::Tab,
CursorKeyMode::default(),
))
.unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
session
.input_event_ch_tx_half
.send(PtyInputEvent::Write(b"After tab".to_vec()))
.unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::Enter,
CursorKeyMode::default(),
))
.unwrap();
tokio::time::sleep(Duration::from_millis(50)).await;
session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::CtrlD,
CursorKeyMode::default(),
))
.unwrap();
let mut output = String::new();
let result = timeout(Duration::from_secs(5), async {
while let Some(event) = session.output_event_receiver_half.recv().await {
match event {
PtyReadWriteOutputEvent::Output(data) => {
output.push_str(&String::from_utf8_lossy(&data));
}
PtyReadWriteOutputEvent::Exit(_) => break,
_ => {}
}
}
})
.await;
assert!(result.is_ok(), "Test timed out. Output: '{output}'");
assert!(
output.contains("Test line"),
"Output should contain 'Test line' but was: '{output}'"
);
assert!(
output.contains("No newline"),
"Output should contain 'No newline' but was: '{output}'"
);
assert!(
output.contains("After tab"),
"Output should contain 'After tab' but was: '{output}'"
);
Ok(())
}
async fn test_raw_escape_sequences() -> miette::Result<()> {
use tokio::time::timeout;
if is_ci::uncached() {
eprintln!("Skipping test_raw_escape_sequences in CI environment");
return Ok(());
}
let temp_dir = std::env::temp_dir();
let mut session = PtyCommandBuilder::new("cat")
.cwd(temp_dir)
.spawn_read_write(PtySize::default())
.unwrap();
let red_text = b"\x1b[31mRed Text\x1b[0m";
session
.input_event_ch_tx_half
.send(PtyInputEvent::Write(red_text.to_vec()))
.unwrap();
session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::Enter,
CursorKeyMode::default(),
))
.unwrap();
tokio::time::sleep(Duration::from_millis(50)).await;
let blue_seq = vec![0x1b, b'[', b'3', b'4', b'm']; session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::RawSequence(blue_seq),
CursorKeyMode::default(),
))
.unwrap();
session
.input_event_ch_tx_half
.send(PtyInputEvent::Write(b"Blue Text".to_vec()))
.unwrap();
let reset_seq = vec![0x1b, b'[', b'0', b'm']; session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::RawSequence(reset_seq),
CursorKeyMode::default(),
))
.unwrap();
session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::Enter,
CursorKeyMode::default(),
))
.unwrap();
tokio::time::sleep(Duration::from_millis(100)).await;
session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::CtrlD,
CursorKeyMode::default(),
))
.unwrap();
let mut output = Vec::new();
let mut events_received = Vec::new();
let mut saw_exit = false;
let result = timeout(Duration::from_secs(10), async {
while let Some(event) = session.output_event_receiver_half.recv().await {
match &event {
PtyReadWriteOutputEvent::Output(data) => {
output.extend_from_slice(data);
events_received.push(format!("Output({} bytes)", data.len()));
}
PtyReadWriteOutputEvent::Exit(status) => {
saw_exit = true;
events_received.push(format!("Exit({status:?})"));
break;
}
other => {
events_received.push(format!("{other:?}"));
}
}
}
})
.await;
let output_str = String::from_utf8_lossy(&output);
assert!(
result.is_ok(),
"Test timed out after 10 seconds. Events received: {events_received:?}, Output so far: '{output_str}'"
);
assert!(
saw_exit,
"Should see exit event. Events received: {events_received:?}, Output: '{output_str}'"
);
assert!(
output_str.contains("Red Text"),
"Output should contain 'Red Text'. Events received: {events_received:?}, Full output was: '{output_str}'"
);
assert!(
output_str.contains("Blue Text"),
"Output should contain 'Blue Text'. Events received: {events_received:?}, Full output was: '{output_str}'"
);
Ok(())
}
#[cfg(not(target_os = "windows"))] #[allow(clippy::too_many_lines)]
async fn test_htop_interactive_with_cursor_modes() -> miette::Result<()> {
use std::process::Command;
use tokio::time::timeout;
println!("Starting htop integration test...");
let htop_check = Command::new("which")
.arg("htop")
.output()
.expect("Failed to check for htop");
assert!(
htop_check.status.success(),
"htop is required for this test but is not installed!\n\
Please install htop:\n\
- Linux: Use your package manager (apt, dnf, pacman, etc.)\n\
- macOS: brew install htop\n\
- Or run: ./bootstrap.sh"
);
println!("htop found, launching with 10-second refresh delay...");
let mut session = timeout(Duration::from_secs(5), async {
PtyCommandBuilder::new("htop")
.args(["--delay", "100"])
.spawn_read_write(PtySize::default())
})
.await
.map_err(|_| miette::miette!("Timeout launching htop"))?
.unwrap();
println!("htop launched successfully, waiting for initialization...");
tokio::time::sleep(Duration::from_millis(1000)).await;
println!("Phase 1: Capturing initial htop display...");
let initial_output = timeout(Duration::from_secs(2), async {
capture_output_snapshot(&mut session, Duration::from_millis(500)).await
})
.await
.map_err(|_| miette::miette!("Timeout capturing initial output"))?
.unwrap_or_else(|_| String::new());
println!("Initial output captured: {} chars", initial_output.len());
assert!(
initial_output.len() > 200,
"htop should display substantial process information. Got {} chars: {}",
initial_output.len(),
initial_output.chars().take(100).collect::<String>()
);
let htop_indicators = initial_output.contains("Tasks:")
|| initial_output.contains("Load average:")
|| initial_output.contains("Memory:")
|| initial_output.contains("PID")
|| initial_output.contains("CPU%")
|| initial_output.len() > 500;
assert!(
htop_indicators,
"htop should display typical process manager content (Tasks, Load average, Memory, PID, CPU%). \
First 200 chars: {}",
initial_output.chars().take(200).collect::<String>()
);
println!("Phase 2: Testing arrow down navigation...");
let send_result =
session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::ArrowDown,
CursorKeyMode::Normal,
));
assert!(send_result.is_ok(), "Failed to send arrow down key");
tokio::time::sleep(Duration::from_millis(300)).await;
let after_arrow = timeout(Duration::from_secs(2), async {
capture_output_snapshot(&mut session, Duration::from_millis(500)).await
})
.await
.map_err(|_| miette::miette!("Timeout capturing post-arrow output"))?
.unwrap_or_else(|_| String::new());
println!("Post-arrow output captured: {} chars", after_arrow.len());
assert!(
after_arrow.len() > 100
&& (after_arrow != initial_output
|| after_arrow.len() != initial_output.len()),
"Arrow down should cause visible UI change in htop. \
Initial: {} chars, After arrow: {} chars",
initial_output.len(),
after_arrow.len()
);
println!("Phase 3: Testing Ctrl+L screen refresh...");
let ctrl_l_result =
session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::CtrlL,
CursorKeyMode::default(),
));
assert!(ctrl_l_result.is_ok(), "Failed to send Ctrl+L");
tokio::time::sleep(Duration::from_millis(300)).await;
let after_ctrl_l = timeout(Duration::from_secs(2), async {
capture_output_snapshot(&mut session, Duration::from_millis(400)).await
})
.await
.map_err(|_| miette::miette!("Timeout capturing post-Ctrl+L output"))?
.unwrap_or_else(|_| String::new());
assert!(
after_ctrl_l.len() > 100,
"Ctrl+L should trigger screen redraw with substantial output. Got {} chars",
after_ctrl_l.len()
);
println!("Phase 4: Testing cursor mode compatibility...");
let mut mode_success_count = 0;
for (mode, mode_name) in [
(CursorKeyMode::Normal, "Normal"),
(CursorKeyMode::Application, "Application"),
] {
let arrow_result = session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(ControlSequence::ArrowUp, mode));
if arrow_result.is_ok() {
mode_success_count += 1;
println!("✓ Arrow key sent successfully in {mode_name} mode");
} else {
println!("✗ Failed to send arrow key in {mode_name} mode");
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert_eq!(
mode_success_count, 2,
"Both cursor modes (Normal and Application) should work. Only {mode_success_count} succeeded"
);
println!("Capturing normal display before F2...");
let before_f2 = timeout(Duration::from_secs(2), async {
capture_output_snapshot(&mut session, Duration::from_millis(300)).await
})
.await
.map_err(|_| miette::miette!("Timeout capturing pre-F2 output"))?
.unwrap_or_else(|_| String::new());
println!("Phase 5: Testing F2 key (Setup menu)...");
let f2_result = session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::F(2),
CursorKeyMode::default(),
));
assert!(f2_result.is_ok(), "Failed to send F2 key");
tokio::time::sleep(Duration::from_millis(400)).await;
let setup_output = timeout(Duration::from_secs(2), async {
capture_output_snapshot(&mut session, Duration::from_millis(500)).await
})
.await
.map_err(|_| miette::miette!("Timeout capturing setup menu output"))?
.unwrap_or_else(|_| String::new());
println!("Setup menu output captured: {} chars", setup_output.len());
let setup_menu_appeared = setup_output != before_f2
&& (setup_output.contains("Setup")
|| setup_output.contains("Meters")
|| setup_output.contains("Display")
|| setup_output.contains("Colors")
|| setup_output.contains("Columns")
|| (!setup_output.is_empty() && setup_output.len() != before_f2.len()));
assert!(
setup_menu_appeared,
"F2 key should open setup menu with different content than normal display.\n\
Before F2: {} chars, Setup menu: {} chars\n\
Setup contains expected keywords: {}\n\
First 200 chars of setup output: {}",
before_f2.len(),
setup_output.len(),
setup_output.contains("Setup")
|| setup_output.contains("Meters")
|| setup_output.contains("Display"),
setup_output.chars().take(200).collect::<String>()
);
println!("Exiting setup menu with Escape...");
let escape_result =
session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::Escape,
CursorKeyMode::default(),
));
assert!(escape_result.is_ok(), "Failed to send Escape key");
tokio::time::sleep(Duration::from_millis(200)).await;
println!("Phase 6: Testing graceful shutdown...");
let quit_result = session
.input_event_ch_tx_half
.send(PtyInputEvent::Write(b"q".to_vec()));
assert!(quit_result.is_ok(), "Failed to send quit command");
let exit_result = timeout(Duration::from_secs(2), async {
while let Some(event) = session.output_event_receiver_half.recv().await {
if let PtyReadWriteOutputEvent::Exit(status) = event {
return Ok(status);
}
}
Err(miette::miette!("No exit event received"))
})
.await;
match exit_result {
Ok(Ok(status)) => {
println!("✓ htop exited cleanly with status: {status:?}");
if !status.success() {
println!(
"⚠ htop exited with non-zero status, but this may be acceptable in test environments"
);
}
}
Ok(Err(_)) | Err(_) => {
println!(
"⚠ htop did not exit gracefully within timeout - applying force termination"
);
let ctrl_c_result =
session
.input_event_ch_tx_half
.send(PtyInputEvent::SendControl(
ControlSequence::CtrlC,
CursorKeyMode::default(),
));
if ctrl_c_result.is_ok() {
tokio::time::sleep(Duration::from_millis(100)).await;
println!("✓ Sent Ctrl+C as fallback termination");
} else {
println!("⚠ Failed to send Ctrl+C fallback - process may be hung");
}
}
}
println!("All phases completed - validating overall test success...");
assert!(
initial_output.len() > 200
|| after_arrow.len() > 200
|| setup_output.len() > 50,
"Expected substantial output from at least one phase. \
Initial: {} chars, After arrow: {} chars, Setup menu: {} chars",
initial_output.len(),
after_arrow.len(),
setup_output.len()
);
println!("✅ htop integration test completed successfully!");
Ok(())
}
async fn capture_output_snapshot(
session: &mut PtyReadWriteSession,
timeout_duration: Duration,
) -> miette::Result<String> {
let mut output = String::new();
let deadline = tokio::time::Instant::now() + timeout_duration;
while tokio::time::Instant::now() < deadline {
tokio::select! {
Some(event) = session.output_event_receiver_half.recv() => {
if let PtyReadWriteOutputEvent::Output(data) = event {
output.push_str(&String::from_utf8_lossy(&data));
}
}
() = tokio::time::sleep_until(deadline) => break,
}
}
Ok(output)
}
#[tokio::test]
async fn test_create_input_handler_task_write() {
let (controller, _controlled) = create_pty_pair(PtySize::default()).unwrap();
let (input_sender, input_receiver) = unbounded_channel();
let (event_sender, _event_receiver) = unbounded_channel();
let handle =
create_controller_input_writer_task(controller, input_receiver, event_sender);
let test_data = b"test input";
input_sender
.send(PtyInputEvent::Write(test_data.to_vec()))
.unwrap();
input_sender.send(PtyInputEvent::Close).unwrap();
tokio::time::sleep(Duration::from_millis(100)).await;
drop(input_sender);
let result = tokio::time::timeout(Duration::from_millis(2000), handle).await;
assert!(result.is_ok(), "Task timed out");
}
#[tokio::test]
async fn test_create_input_handler_task_write_line() {
let (controller, _controlled) = create_pty_pair(PtySize::default()).unwrap();
let (input_sender, input_receiver) = unbounded_channel();
let (event_sender, _event_receiver) = unbounded_channel();
let handle =
create_controller_input_writer_task(controller, input_receiver, event_sender);
input_sender
.send(PtyInputEvent::WriteLine("test line".to_string()))
.unwrap();
input_sender.send(PtyInputEvent::Close).unwrap();
tokio::time::sleep(Duration::from_millis(100)).await;
drop(input_sender);
let result = tokio::time::timeout(Duration::from_millis(2000), handle).await;
assert!(result.is_ok(), "Task timed out");
}
#[tokio::test]
async fn test_create_input_handler_task_control_char() {
let (controller, _controlled) = create_pty_pair(PtySize::default()).unwrap();
let (input_sender, input_receiver) = unbounded_channel();
let (event_sender, _event_receiver) = unbounded_channel();
let handle =
create_controller_input_writer_task(controller, input_receiver, event_sender);
input_sender
.send(PtyInputEvent::SendControl(
ControlSequence::CtrlC,
CursorKeyMode::default(),
))
.unwrap();
input_sender.send(PtyInputEvent::Close).unwrap();
tokio::time::sleep(Duration::from_millis(100)).await;
drop(input_sender);
let result = tokio::time::timeout(Duration::from_millis(2000), handle).await;
assert!(result.is_ok(), "Task timed out");
}
#[tokio::test]
async fn test_create_input_handler_task_resize() {
let (controller, _controlled) = create_pty_pair(PtySize::default()).unwrap();
let (input_sender, input_receiver) = unbounded_channel();
let (event_sender, _event_receiver) = unbounded_channel();
let handle =
create_controller_input_writer_task(controller, input_receiver, event_sender);
let new_size = PtySize {
rows: 40,
cols: 120,
pixel_width: 0,
pixel_height: 0,
};
input_sender.send(PtyInputEvent::Resize(new_size)).unwrap();
input_sender.send(PtyInputEvent::Close).unwrap();
tokio::time::sleep(Duration::from_millis(100)).await;
drop(input_sender);
let result = tokio::time::timeout(Duration::from_millis(2000), handle).await;
assert!(result.is_ok(), "Task timed out");
}
#[tokio::test]
async fn test_create_input_handler_task_flush() {
let (controller, _controlled) = create_pty_pair(PtySize::default()).unwrap();
let (input_sender, input_receiver) = unbounded_channel();
let (event_sender, _event_receiver) = unbounded_channel();
let handle =
create_controller_input_writer_task(controller, input_receiver, event_sender);
input_sender.send(PtyInputEvent::Flush).unwrap();
input_sender.send(PtyInputEvent::Close).unwrap();
tokio::time::sleep(Duration::from_millis(100)).await;
drop(input_sender);
let result = tokio::time::timeout(Duration::from_millis(2000), handle).await;
assert!(result.is_ok(), "Task timed out");
}
#[tokio::test]
async fn test_create_input_handler_task_channel_disconnect() {
let (controller, _controlled) = create_pty_pair(PtySize::default()).unwrap();
let (input_sender, input_receiver) = unbounded_channel();
let (event_sender, _event_receiver) = unbounded_channel();
let handle =
create_controller_input_writer_task(controller, input_receiver, event_sender);
drop(input_sender);
let result = tokio::time::timeout(Duration::from_millis(500), handle).await;
assert!(result.is_ok());
}
}