use std::collections::VecDeque;
use std::io::Write as _;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::thread;
use std::time::{Duration, Instant};
use super::backend::PtySlave;
use super::backend::{Backend, PtyBackend, PtyChild, PtyMaster, PtySize};
use super::{
is_ignorable_process_control_error, poll_pty_process, record_pty_input_metrics,
spawn_pty_reader, store_pty_returncode, terminal_input_relay_worker, write_pty_input,
IdleDetectorCore, NativePtyHandles, PtyError, PtyReadShared, PtyReadState,
};
use running_process_platform_internal::platform::terminal as pty_platform;
pub struct NativePtyProcess {
pub argv: Vec<String>,
pub cwd: Option<String>,
pub env: Option<Vec<(String, String)>>,
pub rows: u16,
pub cols: u16,
pub nice: Option<i32>,
pub handles: Arc<Mutex<Option<NativePtyHandles>>>,
pub reader: Arc<PtyReadShared>,
pub returncode: Arc<Mutex<Option<i32>>>,
pub input_bytes_total: Arc<AtomicUsize>,
pub newline_events_total: Arc<AtomicUsize>,
pub submit_events_total: Arc<AtomicUsize>,
pub echo: Arc<AtomicBool>,
pub idle_detector: Arc<Mutex<Option<Arc<IdleDetectorCore>>>>,
pub output_bytes_total: Arc<AtomicUsize>,
pub control_churn_bytes_total: Arc<AtomicUsize>,
pub reader_worker: Mutex<Option<thread::JoinHandle<()>>>,
pub terminal_input_relay_stop: Arc<AtomicBool>,
pub terminal_input_relay_active: Arc<AtomicBool>,
pub terminal_input_relay_worker: Mutex<Option<thread::JoinHandle<()>>>,
}
pub(super) fn resolved_spawn_cwd(cwd: Option<&str>) -> Option<String> {
cwd.map(str::to_owned).or_else(|| {
std::env::current_dir()
.ok()
.map(|cwd| cwd.to_string_lossy().to_string())
})
}
impl NativePtyProcess {
pub fn new(
argv: Vec<String>,
cwd: Option<String>,
env: Option<Vec<(String, String)>>,
rows: u16,
cols: u16,
nice: Option<i32>,
) -> Result<Self, PtyError> {
if argv.is_empty() {
return Err(PtyError::Other("command cannot be empty".into()));
}
Ok(Self {
argv,
cwd,
env,
rows,
cols,
nice,
handles: Arc::new(Mutex::new(None)),
reader: Arc::new(PtyReadShared {
state: Mutex::new(PtyReadState {
chunks: VecDeque::new(),
closed: false,
}),
condvar: Condvar::new(),
}),
returncode: Arc::new(Mutex::new(None)),
input_bytes_total: Arc::new(AtomicUsize::new(0)),
newline_events_total: Arc::new(AtomicUsize::new(0)),
submit_events_total: Arc::new(AtomicUsize::new(0)),
echo: Arc::new(AtomicBool::new(false)),
idle_detector: Arc::new(Mutex::new(None)),
output_bytes_total: Arc::new(AtomicUsize::new(0)),
control_churn_bytes_total: Arc::new(AtomicUsize::new(0)),
reader_worker: Mutex::new(None),
terminal_input_relay_stop: Arc::new(AtomicBool::new(false)),
terminal_input_relay_active: Arc::new(AtomicBool::new(false)),
terminal_input_relay_worker: Mutex::new(None),
})
}
pub fn mark_reader_closed(&self) {
let mut guard = self.reader.state.lock().expect("pty read mutex poisoned");
guard.closed = true;
self.reader.condvar.notify_all();
}
pub fn store_returncode(&self, code: i32) {
store_pty_returncode(&self.returncode, code);
}
pub fn record_input_metrics(&self, data: &[u8], submit: bool) {
record_pty_input_metrics(
&self.input_bytes_total,
&self.newline_events_total,
&self.submit_events_total,
data,
submit,
);
}
pub fn write_impl(&self, data: &[u8], submit: bool) -> Result<(), PtyError> {
self.record_input_metrics(data, submit);
write_pty_input(&self.handles, data)?;
Ok(())
}
pub fn request_terminal_input_relay_stop(&self) {
self.terminal_input_relay_stop
.store(true, Ordering::Release);
self.terminal_input_relay_active
.store(false, Ordering::Release);
}
pub fn start_terminal_input_relay_impl(&self) -> Result<(), PtyError> {
let mut worker_guard = self
.terminal_input_relay_worker
.lock()
.expect("pty terminal input relay mutex poisoned");
if worker_guard.is_some() && self.terminal_input_relay_active() {
return Ok(());
}
if self
.handles
.lock()
.expect("pty handles mutex poisoned")
.is_none()
{
return Err(PtyError::NotRunning);
}
let Some(input) = pty_platform::TerminalInputSession::new().map_err(PtyError::Io)? else {
self.terminal_input_relay_active
.store(false, Ordering::Release);
return Ok(());
};
self.terminal_input_relay_stop
.store(false, Ordering::Release);
self.terminal_input_relay_active
.store(true, Ordering::Release);
let relay_state = super::TerminalInputRelayState {
handles: Arc::clone(&self.handles),
returncode: Arc::clone(&self.returncode),
input_bytes_total: Arc::clone(&self.input_bytes_total),
newline_events_total: Arc::clone(&self.newline_events_total),
submit_events_total: Arc::clone(&self.submit_events_total),
stop: Arc::clone(&self.terminal_input_relay_stop),
active: Arc::clone(&self.terminal_input_relay_active),
};
*worker_guard = Some(thread::spawn(move || {
terminal_input_relay_worker(input, relay_state);
}));
Ok(())
}
pub fn stop_terminal_input_relay_impl(&self) {
self.request_terminal_input_relay_stop();
if let Some(worker) = self
.terminal_input_relay_worker
.lock()
.expect("pty terminal input relay mutex poisoned")
.take()
{
let _ = worker.join();
}
}
pub fn terminal_input_relay_active(&self) -> bool {
self.terminal_input_relay_active.load(Ordering::Acquire)
}
#[inline(never)]
pub fn close_impl(&self) -> Result<(), PtyError> {
crate::rp_rust_debug_scope!("running_process::NativePtyProcess::close_impl");
self.stop_terminal_input_relay_impl();
let mut guard = self.handles.lock().expect("pty handles mutex poisoned");
let Some(handles) = guard.take() else {
self.mark_reader_closed();
return Ok(());
};
drop(guard);
let NativePtyHandles {
master,
writer,
mut child,
process_guard,
} = handles;
let wait_before_close = pty_platform::wait_before_close_supported();
let mut control_error = None;
if wait_before_close {
if let Err(error) = pty_platform::kill_pty_process_group(master.as_ref()) {
if !is_ignorable_process_control_error(&error) {
control_error = Some(error);
}
}
if let Err(error) = child.kill() {
if !is_ignorable_process_control_error(&error) && control_error.is_none() {
control_error = Some(error);
}
}
}
drop(process_guard);
let reap_deadline = Instant::now() + Duration::from_secs(2);
let code = loop {
match child.try_wait() {
Ok(Some(status)) => break status as i32,
Ok(None) if Instant::now() < reap_deadline => {
thread::sleep(Duration::from_millis(10));
}
Ok(None) if wait_before_close => break -9,
Ok(None) => {
if let Err(error) = child.kill() {
if !is_ignorable_process_control_error(&error) && control_error.is_none() {
control_error = Some(error);
}
}
let kill_deadline = Instant::now() + Duration::from_secs(2);
break loop {
match child.try_wait() {
Ok(Some(status)) => break status as i32,
Ok(None) if Instant::now() < kill_deadline => {
thread::sleep(Duration::from_millis(10));
}
_ => break -9,
}
};
}
Err(error) => {
if control_error.is_none() {
control_error = Some(error);
}
break -9;
}
}
};
drop(writer);
let reader_worker = self
.reader_worker
.lock()
.expect("pty reader worker mutex poisoned")
.take();
let (teardown_tx, teardown_rx) = std::sync::mpsc::channel();
thread::spawn(move || {
drop(master);
drop(child);
if let Some(worker) = reader_worker {
let _ = worker.join();
}
let _ = teardown_tx.send(());
});
let _ = teardown_rx.recv_timeout(Duration::from_secs(2));
self.store_returncode(code);
self.mark_reader_closed();
control_error.map_or(Ok(()), |error| Err(PtyError::Io(error)))
}
#[inline(never)]
pub fn close_nonblocking(&self) {
crate::rp_rust_debug_scope!("running_process::NativePtyProcess::close_nonblocking");
self.request_terminal_input_relay_stop();
let Ok(mut guard) = self.handles.lock() else {
return;
};
let Some(handles) = guard.take() else {
self.mark_reader_closed();
return;
};
drop(guard);
let NativePtyHandles {
master,
writer,
mut child,
process_guard,
} = handles;
let _ = child.kill();
drop(writer);
if pty_platform::wait_before_close_supported() {
drop(master);
drop(child);
drop(process_guard);
} else {
thread::spawn(move || {
drop(master);
drop(child);
drop(process_guard);
});
}
self.mark_reader_closed();
}
pub fn start_impl(&self) -> Result<(), PtyError> {
crate::rp_rust_debug_scope!("running_process::NativePtyProcess::start");
let mut guard = self.handles.lock().expect("pty handles mutex poisoned");
if guard.is_some() {
return Err(PtyError::AlreadyStarted);
}
let spawn_context = pty_platform::before_pty_spawn();
let (mut master, slave) = Backend::openpty(PtySize {
rows: self.rows,
cols: self.cols,
pixel_width: 0,
pixel_height: 0,
})
.map_err(|e| PtyError::Spawn(e.to_string()))?;
let argv: Vec<std::ffi::OsString> =
self.argv.iter().map(std::ffi::OsString::from).collect();
let cwd = resolved_spawn_cwd(self.cwd.as_deref());
let env: Option<Vec<(std::ffi::OsString, std::ffi::OsString)>> =
self.env.as_ref().map(|e| {
e.iter()
.map(|(k, v)| (std::ffi::OsString::from(k), std::ffi::OsString::from(v)))
.collect()
});
let reader = master
.try_clone_reader()
.map_err(|e| PtyError::Spawn(e.to_string()))?;
let writer = master
.take_writer()
.map_err(|e| PtyError::Spawn(e.to_string()))?;
let cwd_path = cwd.as_deref().map(std::path::Path::new);
let child = slave
.spawn(&argv, cwd_path, env.as_deref())
.map_err(|e| PtyError::Spawn(e.to_string()))?;
let process_guard = pty_platform::prepare_pty_child(spawn_context, &child, self.nice)
.map_err(PtyError::Io)?;
let shared = Arc::clone(&self.reader);
let echo = Arc::clone(&self.echo);
let idle_detector = Arc::clone(&self.idle_detector);
let output_bytes = Arc::clone(&self.output_bytes_total);
let churn_bytes = Arc::clone(&self.control_churn_bytes_total);
let reader_worker = thread::spawn(move || {
spawn_pty_reader(
reader,
shared,
echo,
idle_detector,
output_bytes,
churn_bytes,
);
});
*self
.reader_worker
.lock()
.expect("pty reader worker mutex poisoned") = Some(reader_worker);
*guard = Some(NativePtyHandles {
master: Box::new(master) as Box<dyn PtyMaster>,
writer: Arc::new(Mutex::new(writer)),
child: Box::new(child) as Box<dyn PtyChild>,
process_guard,
});
Ok(())
}
pub fn respond_to_queries_impl(&self, data: &[u8]) -> Result<(), PtyError> {
let responses = pty_platform::query_responses(data);
if responses.is_empty() {
return Ok(());
}
let writer = {
let guard = self.handles.lock().expect("pty handles mutex poisoned");
let handles = guard.as_ref().ok_or(PtyError::NotRunning)?;
Arc::clone(&handles.writer)
};
let mut writer = writer.lock().expect("pty writer mutex poisoned");
for response in responses {
writer.write_all(&response).map_err(PtyError::Io)?;
}
writer.flush().map_err(PtyError::Io)
}
pub fn resize_impl(&self, rows: u16, cols: u16) -> Result<(), PtyError> {
crate::rp_rust_debug_scope!("running_process::NativePtyProcess::resize");
let guard = self.handles.lock().expect("pty handles mutex poisoned");
if let Some(handles) = guard.as_ref() {
pty_platform::resize_pty(
handles.master.as_ref(),
PtySize {
rows,
cols,
pixel_width: 0,
pixel_height: 0,
},
)
.map_err(|error| PtyError::Other(error.to_string()))?;
}
Ok(())
}
pub fn send_interrupt_impl(&self) -> Result<(), PtyError> {
crate::rp_rust_debug_scope!("running_process::NativePtyProcess::send_interrupt");
let (target, writer) = {
let guard = self.handles.lock().expect("pty handles mutex poisoned");
let handles = guard.as_ref().ok_or(PtyError::NotRunning)?;
(
handles.master.interrupt_target().map_err(PtyError::Io)?,
Arc::clone(&handles.writer),
)
};
let wrote_input =
pty_platform::send_pty_interrupt(target, &writer).map_err(PtyError::Io)?;
if wrote_input {
self.record_input_metrics(&[0x03], false);
}
Ok(())
}
pub fn wait_impl(&self, timeout: Option<f64>) -> Result<i32, PtyError> {
crate::rp_rust_debug_scope!("running_process::NativePtyProcess::wait");
if let Some(code) = *self
.returncode
.lock()
.expect("pty returncode mutex poisoned")
{
return Ok(code);
}
let start = Instant::now();
loop {
if let Some(code) = poll_pty_process(&self.handles, &self.returncode)? {
return Ok(code);
}
if timeout.is_some_and(|limit| start.elapsed() >= Duration::from_secs_f64(limit)) {
return Err(PtyError::Timeout);
}
thread::sleep(Duration::from_millis(10));
}
}
pub fn terminate_impl(&self) -> Result<(), PtyError> {
crate::rp_rust_debug_scope!("running_process::NativePtyProcess::terminate");
let should_close = {
let mut guard = self.handles.lock().expect("pty handles mutex poisoned");
let handles = guard.as_mut().ok_or(PtyError::NotRunning)?;
let pid = handles.child.pid();
if pid == 0 {
return Err(PtyError::NotRunning);
}
pty_platform::terminate_pty_child(pid).map_err(PtyError::Io)?
};
if should_close {
self.close_impl()?;
}
Ok(())
}
pub fn kill_impl(&self) -> Result<(), PtyError> {
crate::rp_rust_debug_scope!("running_process::NativePtyProcess::kill");
if self
.handles
.lock()
.expect("pty handles mutex poisoned")
.is_none()
{
return Err(PtyError::NotRunning);
}
self.close_impl()
}
pub fn terminate_tree_impl(&self) -> Result<(), PtyError> {
crate::rp_rust_debug_scope!("running_process::NativePtyProcess::terminate_tree");
let Some(pid) = self.pid()? else {
return if self
.returncode
.lock()
.expect("pty returncode mutex poisoned")
.is_some()
{
Ok(())
} else {
Err(PtyError::NotRunning)
};
};
if pty_platform::signal_pty_tree(pid, false).map_err(PtyError::Io)? {
self.close_impl()?;
}
Ok(())
}
pub fn kill_tree_impl(&self) -> Result<(), PtyError> {
crate::rp_rust_debug_scope!("running_process::NativePtyProcess::kill_tree");
let Some(pid) = self.pid()? else {
return if self
.returncode
.lock()
.expect("pty returncode mutex poisoned")
.is_some()
{
Ok(())
} else {
Err(PtyError::NotRunning)
};
};
if pty_platform::signal_pty_tree(pid, true).map_err(PtyError::Io)? {
self.close_impl()?;
}
Ok(())
}
pub fn pid(&self) -> Result<Option<u32>, PtyError> {
let guard = self.handles.lock().expect("pty handles mutex poisoned");
if let Some(handles) = guard.as_ref() {
return Ok(pty_platform::preferred_pty_pid(
handles.master.as_ref(),
handles.child.as_ref(),
));
}
Ok(None)
}
pub fn read_chunk_impl(&self, timeout: Option<f64>) -> Result<Option<Vec<u8>>, PtyError> {
let deadline = timeout.map(|secs| Instant::now() + Duration::from_secs_f64(secs));
let mut guard = self.reader.state.lock().expect("pty read mutex poisoned");
loop {
if let Some(chunk) = guard.chunks.pop_front() {
return Ok(Some(chunk));
}
if guard.closed {
return Err(PtyError::Other("Pseudo-terminal stream is closed".into()));
}
match deadline {
Some(deadline) => {
let now = Instant::now();
if now >= deadline {
return Ok(None); }
let wait = deadline.saturating_duration_since(now);
let result = self
.reader
.condvar
.wait_timeout(guard, wait)
.expect("pty read mutex poisoned");
guard = result.0;
}
None => {
guard = self
.reader
.condvar
.wait(guard)
.expect("pty read mutex poisoned");
}
}
}
}
pub fn wait_for_reader_closed_impl(&self, timeout: Option<f64>) -> bool {
let deadline = timeout.map(|secs| Instant::now() + Duration::from_secs_f64(secs));
let mut guard = self.reader.state.lock().expect("pty read mutex poisoned");
loop {
if guard.closed {
return true;
}
match deadline {
Some(deadline) => {
let now = Instant::now();
if now >= deadline {
return false;
}
let wait = deadline.saturating_duration_since(now);
let result = self
.reader
.condvar
.wait_timeout(guard, wait)
.expect("pty read mutex poisoned");
guard = result.0;
}
None => {
guard = self
.reader
.condvar
.wait(guard)
.expect("pty read mutex poisoned");
}
}
}
}
pub fn wait_and_drain_impl(
&self,
timeout: Option<f64>,
drain_timeout: f64,
) -> Result<i32, PtyError> {
let code = self.wait_impl(timeout)?;
let deadline = Instant::now() + Duration::from_secs_f64(drain_timeout.max(0.0));
let mut guard = self.reader.state.lock().expect("pty read mutex poisoned");
while !guard.closed {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
break;
}
let result = self
.reader
.condvar
.wait_timeout(guard, remaining)
.expect("pty read mutex poisoned");
guard = result.0;
}
Ok(code)
}
pub fn set_echo(&self, enabled: bool) {
self.echo.store(enabled, Ordering::Release);
}
pub fn echo_enabled(&self) -> bool {
self.echo.load(Ordering::Acquire)
}
pub fn attach_idle_detector(&self, detector: &Arc<IdleDetectorCore>) {
let mut guard = self
.idle_detector
.lock()
.expect("idle detector mutex poisoned");
*guard = Some(Arc::clone(detector));
}
pub fn detach_idle_detector(&self) {
let mut guard = self
.idle_detector
.lock()
.expect("idle detector mutex poisoned");
*guard = None;
}
pub fn pty_input_bytes_total(&self) -> usize {
self.input_bytes_total.load(Ordering::Acquire)
}
pub fn pty_newline_events_total(&self) -> usize {
self.newline_events_total.load(Ordering::Acquire)
}
pub fn pty_submit_events_total(&self) -> usize {
self.submit_events_total.load(Ordering::Acquire)
}
pub fn pty_output_bytes_total(&self) -> usize {
self.output_bytes_total.load(Ordering::Acquire)
}
pub fn pty_control_churn_bytes_total(&self) -> usize {
self.control_churn_bytes_total.load(Ordering::Acquire)
}
}
#[derive(Debug, Clone, Copy)]
pub struct InteractivePtyOptions {
pub echo_output: bool,
pub relay_terminal_input: bool,
pub respond_to_queries: bool,
}
impl Default for InteractivePtyOptions {
fn default() -> Self {
Self {
echo_output: true,
relay_terminal_input: true,
respond_to_queries: true,
}
}
}
#[derive(Debug, Default)]
pub struct InteractivePtyPumpResult {
pub chunks: Vec<Vec<u8>>,
pub stream_closed: bool,
}
pub struct InteractivePtySession {
process: NativePtyProcess,
options: InteractivePtyOptions,
}
impl InteractivePtySession {
pub fn new(process: NativePtyProcess) -> Self {
Self::with_options(process, InteractivePtyOptions::default())
}
pub fn with_options(process: NativePtyProcess, options: InteractivePtyOptions) -> Self {
Self { process, options }
}
pub fn process(&self) -> &NativePtyProcess {
&self.process
}
pub fn start(&self) -> Result<(), PtyError> {
self.process.set_echo(self.options.echo_output);
self.process.start_impl()?;
if self.options.relay_terminal_input {
self.process.start_terminal_input_relay_impl()?;
}
Ok(())
}
pub fn pump_output(
&self,
timeout: Option<f64>,
consume_all: bool,
) -> Result<InteractivePtyPumpResult, PtyError> {
let mut pumped = InteractivePtyPumpResult::default();
let mut next_timeout = timeout;
loop {
match self.process.read_chunk_impl(next_timeout) {
Ok(Some(chunk)) => {
if self.options.respond_to_queries {
self.process.respond_to_queries_impl(&chunk)?;
}
pumped.chunks.push(chunk);
if !consume_all {
break;
}
next_timeout = Some(0.0);
}
Ok(None) => break,
Err(PtyError::Other(message)) if message == "Pseudo-terminal stream is closed" => {
pumped.stream_closed = true;
break;
}
Err(err) => return Err(err),
}
}
Ok(pumped)
}
pub fn resize(&self, rows: u16, cols: u16) -> Result<(), PtyError> {
self.process.resize_impl(rows, cols)
}
pub fn send_interrupt(&self) -> Result<(), PtyError> {
self.process.send_interrupt_impl()
}
pub fn wait(&self, timeout: Option<f64>) -> Result<i32, PtyError> {
self.process.wait_impl(timeout)
}
pub fn wait_and_drain(
&self,
timeout: Option<f64>,
drain_timeout: f64,
) -> Result<i32, PtyError> {
self.process.wait_and_drain_impl(timeout, drain_timeout)
}
pub fn terminate(&self) -> Result<(), PtyError> {
self.process.terminate_impl()
}
pub fn kill(&self) -> Result<(), PtyError> {
self.process.kill_impl()
}
pub fn close(&self) -> Result<(), PtyError> {
self.process.close_impl()
}
}
impl Drop for NativePtyProcess {
fn drop(&mut self) {
self.close_nonblocking();
}
}
#[cfg(test)]
#[path = "../tests/native_pty_process_coverage.rs"]
mod coverage_tests;