use crate::{Controlled, ControlledChild, Controller, ControllerReader, OscBuffer,
PtyCommandBuilder, PtyConfig, PtyReadOnlyOutputEvent, PtyReadOnlySession,
READ_BUFFER_SIZE,
pty_common_io::{create_pty_pair, spawn_command_in_pty}};
use miette::IntoDiagnostic;
use std::io::Read;
impl PtyCommandBuilder {
pub fn spawn_read_only(
self,
arg_config: impl Into<PtyConfig>,
) -> miette::Result<PtyReadOnlySession> {
let pty_config = arg_config.into();
let (
output_evt_ch_tx_half,
output_evt_ch_rx_half,
) = tokio::sync::mpsc::unbounded_channel();
let session_completion_handle = tokio::spawn(async move {
let command = self.build()?;
let (controller, controlled): (Controller, Controlled) =
create_pty_pair(pty_config.get_pty_size())?;
let controlled_child: ControlledChild =
spawn_command_in_pty(&controlled, command)?;
let output_reader_task_handle = {
let controller_reader = controller
.try_clone_reader()
.map_err(|e| miette::miette!("Failed to clone pty reader: {}", e))?;
spawn_blocking_controller_output_reader_task(
controller_reader,
output_evt_ch_tx_half.clone(),
pty_config,
)
};
let child_proc_exit_code = spawn_child_process_waiter(
controlled_child,
output_evt_ch_tx_half.clone(),
)
.await
.into_diagnostic()??;
drop(controlled); drop(controller);
output_reader_task_handle.await.into_diagnostic()??;
Ok(portable_pty::ExitStatus::with_exit_code(
child_proc_exit_code,
))
});
Ok(PtyReadOnlySession {
output_evt_ch_rx_half,
pinned_boxed_session_completion_handle: Box::pin(session_completion_handle),
})
}
}
#[must_use]
fn spawn_child_process_waiter(
mut controlled_child: ControlledChild,
output_evt_ch_tx_half: tokio::sync::mpsc::UnboundedSender<PtyReadOnlyOutputEvent>,
) -> tokio::task::JoinHandle<miette::Result<u32>> {
tokio::task::spawn_blocking(move || -> miette::Result<u32> {
let status = controlled_child.wait().into_diagnostic()?;
let exit_code = status.exit_code();
let _unused = output_evt_ch_tx_half.send(PtyReadOnlyOutputEvent::Exit(status));
Ok(exit_code)
})
}
#[must_use]
pub fn spawn_blocking_controller_output_reader_task(
mut controller_reader: ControllerReader,
output_event_ch_tx_half: tokio::sync::mpsc::UnboundedSender<PtyReadOnlyOutputEvent>,
arg_config: impl Into<PtyConfig>,
) -> tokio::task::JoinHandle<miette::Result<()>> {
let pty_config: PtyConfig = arg_config.into();
tokio::task::spawn_blocking(move || -> miette::Result<()> {
let mut read_buffer = [0u8; READ_BUFFER_SIZE];
let mut osc_buffer = if pty_config.is_osc_capture_enabled() {
Some(OscBuffer::new())
} else {
None
};
loop {
match controller_reader.read(&mut read_buffer) {
Ok(0) | Err(_) => break, Ok(n) => {
let data = &read_buffer[..n];
if pty_config.is_output_capture_enabled() {
let _unused = output_event_ch_tx_half
.send(PtyReadOnlyOutputEvent::Output(data.to_vec()));
}
if let Some(ref mut osc_buf) = osc_buffer {
for event in osc_buf.append_and_extract(data, n) {
let _unused = output_event_ch_tx_half
.send(PtyReadOnlyOutputEvent::Osc(event));
}
}
}
}
}
drop(controller_reader);
Ok(())
})
}
#[cfg(test)]
mod tests {
use crate::{OscEvent, PtyCommandBuilder, PtyConfigOption, PtyReadOnlyOutputEvent,
PtyReadOnlySession,
pty_read_only::spawn_blocking_controller_output_reader_task};
use miette::IntoDiagnostic;
use tokio::{sync::mpsc::unbounded_channel,
time::{Duration, timeout}};
async fn collect_events_with_timeout(
mut session: PtyReadOnlySession,
max_duration: Duration,
) -> miette::Result<(Vec<PtyReadOnlyOutputEvent>, portable_pty::ExitStatus)> {
let mut events = Vec::new();
let result = timeout(max_duration, async move {
loop {
tokio::select! {
result = &mut session.pinned_boxed_session_completion_handle => {
let status = result.into_diagnostic()??;
while let Ok(event) = session.output_evt_ch_rx_half.try_recv() {
events.push(event);
}
return Ok::<_, miette::Error>((events, status));
}
Some(event) = session.output_evt_ch_rx_half.recv() => {
events.push(event);
}
}
}
})
.await;
match result {
Ok(Ok(data)) => Ok(data),
Ok(Err(e)) => Err(e),
Err(_) => Err(miette::miette!("Test timed out")),
}
}
#[tokio::test]
async fn test_simple_echo_command() -> miette::Result<()> {
let temp_dir = std::env::temp_dir();
let session = PtyCommandBuilder::new("echo")
.args(["Hello, PTY!"])
.cwd(temp_dir)
.spawn_read_only(PtyConfigOption::Output)?;
let (events, status) =
collect_events_with_timeout(session, Duration::from_secs(5)).await?;
assert!(status.success());
let output_events: Vec<_> = events
.iter()
.filter_map(|e| match e {
PtyReadOnlyOutputEvent::Output(data) => {
Some(String::from_utf8_lossy(data).to_string())
}
_ => None,
})
.collect();
assert!(!output_events.is_empty());
let combined_output = output_events.join("");
assert!(combined_output.contains("Hello, PTY!"));
Ok(())
}
#[tokio::test]
async fn test_osc_sequence_with_printf() -> miette::Result<()> {
let temp_dir = std::env::temp_dir();
let session = PtyCommandBuilder::new("bash")
.args(["-c", r"printf '\033]9;4;1;50\033\\'"])
.cwd(temp_dir)
.spawn_read_only(PtyConfigOption::Osc)?;
let (events, status) =
collect_events_with_timeout(session, Duration::from_secs(5)).await?;
assert!(status.success());
let osc_events: Vec<_> = events
.iter()
.filter_map(|e| match e {
PtyReadOnlyOutputEvent::Osc(osc) => Some(osc.clone()),
_ => None,
})
.collect();
assert_eq!(osc_events, vec![OscEvent::ProgressUpdate(50)]);
Ok(())
}
#[tokio::test]
async fn test_multiple_osc_sequences() -> miette::Result<()> {
let temp_dir = std::env::temp_dir();
let session = PtyCommandBuilder::new("bash")
.args([
"-c",
r"printf '\033]9;4;1;25\033\\\033]9;4;1;50\033\\\033]9;4;0;0\033\\'",
])
.cwd(temp_dir)
.spawn_read_only(PtyConfigOption::Osc)?;
let (events, status) =
collect_events_with_timeout(session, Duration::from_secs(5)).await?;
assert!(status.success());
let osc_events: Vec<_> = events
.iter()
.filter_map(|e| match e {
PtyReadOnlyOutputEvent::Osc(osc) => Some(osc.clone()),
_ => None,
})
.collect();
assert_eq!(
osc_events,
vec![
OscEvent::ProgressUpdate(25),
OscEvent::ProgressUpdate(50),
OscEvent::ProgressCleared,
]
);
Ok(())
}
#[tokio::test]
async fn test_osc_with_mixed_output() -> miette::Result<()> {
if cfg!(target_os = "windows") {
return Ok(()); }
let temp_dir = std::env::temp_dir();
let config = PtyConfigOption::Osc + PtyConfigOption::Output;
let session = PtyCommandBuilder::new("bash")
.args([
"-c",
r"echo 'Starting...'; printf '\033]9;4;1;50\033\\'; echo 'Done!'",
])
.cwd(temp_dir)
.spawn_read_only(config)?;
let (events, status) =
collect_events_with_timeout(session, Duration::from_secs(5)).await?;
assert!(status.success());
let has_output = events
.iter()
.any(|e| matches!(e, PtyReadOnlyOutputEvent::Output(_)));
assert!(has_output);
let osc_events: Vec<_> = events
.iter()
.filter_map(|e| match e {
PtyReadOnlyOutputEvent::Osc(osc) => Some(osc.clone()),
_ => None,
})
.collect();
if !osc_events.is_empty() {
assert_eq!(osc_events, vec![OscEvent::ProgressUpdate(50)]);
}
Ok(())
}
#[tokio::test]
async fn test_split_osc_sequence_simulation() -> miette::Result<()> {
if cfg!(target_os = "windows") {
return Ok(()); }
let temp_dir = std::env::temp_dir();
let session = PtyCommandBuilder::new("bash")
.args(["-c", r"printf '\033]9;4;1;'; sleep 0.01; printf '75\033\\'"])
.cwd(temp_dir)
.spawn_read_only(PtyConfigOption::Osc)?;
let (events, status) =
collect_events_with_timeout(session, Duration::from_secs(5)).await?;
assert!(status.success());
let osc_events: Vec<_> = events
.iter()
.filter_map(|e| match e {
PtyReadOnlyOutputEvent::Osc(osc) => Some(osc.clone()),
_ => None,
})
.collect();
if !osc_events.is_empty() {
assert_eq!(osc_events, vec![OscEvent::ProgressUpdate(75)]);
}
Ok(())
}
#[tokio::test]
async fn test_all_osc_event_types() -> miette::Result<()> {
let temp_dir = std::env::temp_dir();
let sequences = [
(r"printf '\033]9;4;0;0\033\\'", OscEvent::ProgressCleared),
(
r"printf '\033]9;4;1;42\033\\'",
OscEvent::ProgressUpdate(42),
),
(r"printf '\033]9;4;2;0\033\\'", OscEvent::BuildError),
(
r"printf '\033]9;4;3;0\033\\'",
OscEvent::IndeterminateProgress,
),
];
for (bash_cmd, expected) in sequences {
let session = PtyCommandBuilder::new("bash")
.args(["-c", bash_cmd])
.cwd(temp_dir.clone())
.spawn_read_only(PtyConfigOption::Osc)?;
let (events, status) =
collect_events_with_timeout(session, Duration::from_secs(5)).await?;
assert!(status.success());
let osc_events: Vec<_> = events
.iter()
.filter_map(|e| match e {
PtyReadOnlyOutputEvent::Osc(osc) => Some(osc.clone()),
_ => None,
})
.collect();
assert_eq!(osc_events, vec![expected]);
}
Ok(())
}
#[tokio::test]
async fn test_command_failure() -> miette::Result<()> {
let temp_dir = std::env::temp_dir();
let session = PtyCommandBuilder::new("false")
.cwd(temp_dir)
.spawn_read_only(PtyConfigOption::NoCaptureOutput)?;
let (events, status) =
collect_events_with_timeout(session, Duration::from_secs(5)).await?;
assert!(!status.success());
let has_exit = events
.iter()
.any(|e| matches!(e, PtyReadOnlyOutputEvent::Exit(_)));
assert!(has_exit);
Ok(())
}
#[tokio::test]
async fn test_no_capture_option() -> miette::Result<()> {
let temp_dir = std::env::temp_dir();
let session = PtyCommandBuilder::new("echo")
.args(["test"])
.cwd(temp_dir)
.spawn_read_only(PtyConfigOption::NoCaptureOutput)?;
let (events, status) =
collect_events_with_timeout(session, Duration::from_secs(5)).await?;
assert!(status.success());
let output_events: Vec<_> = events
.iter()
.filter(|e| matches!(e, PtyReadOnlyOutputEvent::Output(_)))
.collect();
assert!(output_events.is_empty());
Ok(())
}
#[tokio::test]
async fn test_create_reader_task_no_capture() {
let (event_sender, mut event_receiver) = unbounded_channel();
let mock_data = b"test data";
let reader = Box::new(std::io::Cursor::new(mock_data.to_vec()));
let handle = spawn_blocking_controller_output_reader_task(
reader,
event_sender,
PtyConfigOption::NoCaptureOutput,
);
let result = tokio::time::timeout(Duration::from_millis(100), handle).await;
assert!(result.is_ok());
assert!(event_receiver.try_recv().is_err());
}
#[tokio::test]
async fn test_create_reader_task_with_output_capture() {
let (event_sender, mut event_receiver) = unbounded_channel();
let mock_data = b"test data";
let reader = Box::new(std::io::Cursor::new(mock_data.to_vec()));
let handle = spawn_blocking_controller_output_reader_task(
reader,
event_sender,
PtyConfigOption::Output,
);
let result = tokio::time::timeout(Duration::from_millis(100), handle).await;
assert!(result.is_ok());
if let Ok(event) = event_receiver.try_recv() {
match event {
PtyReadOnlyOutputEvent::Output(data) => assert_eq!(data, mock_data),
_ => panic!("Expected Output event"),
}
}
}
#[tokio::test]
async fn test_create_reader_task_with_osc_capture() {
let (event_sender, mut event_receiver) = unbounded_channel();
let mock_data = b"\x1b]9;4;1;50\x1b\\";
let reader = Box::new(std::io::Cursor::new(mock_data.to_vec()));
let handle = spawn_blocking_controller_output_reader_task(
reader,
event_sender,
PtyConfigOption::Osc,
);
let result = tokio::time::timeout(Duration::from_millis(100), handle).await;
assert!(result.is_ok());
let mut events = Vec::new();
while let Ok(event) = event_receiver.try_recv() {
events.push(event);
}
assert!(
!events.is_empty(),
"Should have received at least one event"
);
let osc_events: Vec<_> = events
.iter()
.filter_map(|e| match e {
PtyReadOnlyOutputEvent::Osc(osc) => Some(osc),
_ => None,
})
.collect();
assert!(
!osc_events.is_empty(),
"Should have received at least one OSC event"
);
let has_correct_event = osc_events
.iter()
.any(|osc| matches!(osc, crate::OscEvent::ProgressUpdate(50)));
assert!(
has_correct_event,
"Expected OSC progress update event with 50%"
);
}
#[tokio::test]
async fn test_create_reader_task_with_osc_only_capture() {
let (event_sender, mut event_receiver) = unbounded_channel();
let mock_data = b"\x1b]9;4;1;75\x1b\\";
let reader = Box::new(std::io::Cursor::new(mock_data.to_vec()));
let config = PtyConfigOption::Osc
+ PtyConfigOption::NoCaptureOutput
+ PtyConfigOption::Osc;
let handle =
spawn_blocking_controller_output_reader_task(reader, event_sender, config);
let result = tokio::time::timeout(Duration::from_millis(100), handle).await;
assert!(result.is_ok());
let mut events = Vec::new();
while let Ok(event) = event_receiver.try_recv() {
events.push(event);
}
assert!(
!events.is_empty(),
"Should have received at least one event"
);
let osc_events: Vec<_> = events
.iter()
.filter_map(|e| match e {
PtyReadOnlyOutputEvent::Osc(osc) => Some(osc),
_ => None,
})
.collect();
let output_events: Vec<_> = events
.iter()
.filter_map(|e| match e {
PtyReadOnlyOutputEvent::Output(_) => Some(()),
_ => None,
})
.collect();
assert!(!osc_events.is_empty(), "Should have received OSC events");
assert!(
output_events.is_empty(),
"Should NOT have received output events (OSC-only capture)"
);
let has_correct_event = osc_events
.iter()
.any(|osc| matches!(osc, crate::OscEvent::ProgressUpdate(75)));
assert!(
has_correct_event,
"Expected OSC progress update event with 75%"
);
}
#[tokio::test]
async fn test_create_reader_task_with_both_osc_and_output_capture() {
let (event_sender, mut event_receiver) = unbounded_channel();
let mock_data = b"\x1b]9;4;1;25\x1b\\";
let reader = Box::new(std::io::Cursor::new(mock_data.to_vec()));
let config = PtyConfigOption::Osc + PtyConfigOption::Output;
let handle =
spawn_blocking_controller_output_reader_task(reader, event_sender, config);
let result = tokio::time::timeout(Duration::from_millis(100), handle).await;
assert!(result.is_ok());
let mut events = Vec::new();
while let Ok(event) = event_receiver.try_recv() {
events.push(event);
}
assert!(
!events.is_empty(),
"Should have received at least one event"
);
let osc_events: Vec<_> = events
.iter()
.filter_map(|e| match e {
PtyReadOnlyOutputEvent::Osc(osc) => Some(osc),
_ => None,
})
.collect();
let output_events: Vec<_> = events
.iter()
.filter_map(|e| match e {
PtyReadOnlyOutputEvent::Output(_) => Some(()),
_ => None,
})
.collect();
assert!(!osc_events.is_empty(), "Should have received OSC events");
assert!(
!output_events.is_empty(),
"Should have received output events (both capture enabled)"
);
let has_correct_event = osc_events
.iter()
.any(|osc| matches!(osc, crate::OscEvent::ProgressUpdate(25)));
assert!(
has_correct_event,
"Expected OSC progress update event with 25%"
);
}
}