use std::io;
use std::os::windows::io::{AsHandle, OwnedHandle};
#[cfg(test)]
use std::sync::atomic::{AtomicU8, Ordering};
#[cfg(any(feature = "blocking", feature = "tokio"))]
use std::sync::mpsc;
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::thread;
#[cfg(any(feature = "blocking", feature = "tokio"))]
use crate::backend::BackendKind;
use crate::backend::{ConPtyBackend, HPCON, PSEUDOCONSOLE_INHERIT_CURSOR};
use crate::size::Size;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ReaderState {
Open,
Drained,
Closed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CloseState {
NotRequested,
Requested,
Done,
}
#[derive(Debug)]
struct State {
size: Size,
released: bool,
release_failed: bool,
reader: ReaderState,
close: CloseState,
#[cfg(test)]
drop_observer: Option<Arc<AtomicU8>>,
}
impl State {
const fn initial(size: Size) -> Self {
Self {
size,
released: false,
release_failed: false,
reader: ReaderState::Open,
close: CloseState::NotRequested,
#[cfg(test)]
drop_observer: None,
}
}
}
#[derive(Debug)]
pub(crate) struct ConsoleShared {
backend: ConPtyBackend,
hpcon: HPCON,
state: Mutex<State>,
}
#[derive(Debug)]
pub(super) struct SpawnCapability<'a> {
hpcon: HPCON,
_state: MutexGuard<'a, State>,
}
unsafe impl windows_spawn::AsPseudoConsole for SpawnCapability<'_> {
fn raw_pseudoconsole(&self) -> isize {
self.hpcon
}
}
impl ConsoleShared {
fn lock(&self) -> MutexGuard<'_, State> {
self.state.lock().unwrap_or_else(PoisonError::into_inner)
}
#[cfg(test)]
fn observe_drop_mode(&self, observer: Arc<AtomicU8>) {
self.lock().drop_observer = Some(observer);
}
pub(crate) fn release_after_spawn(&self) -> io::Result<bool> {
let mut state = self.lock();
if state.close == CloseState::Done {
return Ok(false);
}
if state.released {
return Ok(true);
}
match unsafe { self.backend.release(self.hpcon) } {
None => Ok(false),
Some(Ok(())) => {
state.released = true;
Ok(true)
},
Some(Err(err)) => {
state.release_failed = true;
Err(err)
},
}
}
pub(crate) fn notify_reader_eof(&self) {
let mut state = self.lock();
if state.reader == ReaderState::Open {
state.reader = ReaderState::Drained;
}
if state.close != CloseState::Done {
state.close = CloseState::Requested;
self.execute_close(state);
}
}
pub(crate) fn notify_reader_closed(&self) {
let mut state = self.lock();
state.reader = ReaderState::Closed;
self.close_if_due(state);
}
pub(crate) fn request_close(&self) {
let mut state = self.lock();
if state.close == CloseState::Done {
return;
}
state.close = CloseState::Requested;
if state.reader == ReaderState::Open && state.released {
return;
}
self.execute_close(state);
}
#[cfg(any(feature = "blocking", feature = "tokio"))]
pub(crate) fn request_close_detached(self: &Arc<Self>) {
self.request_close_detached_inner(true);
}
#[cfg(all(test, any(feature = "blocking", feature = "tokio")))]
fn request_close_detached_with_worker_spawn_failure(self: &Arc<Self>) {
self.request_close_detached_inner(false);
}
#[cfg(any(feature = "blocking", feature = "tokio"))]
fn request_close_detached_inner(self: &Arc<Self>, spawn_worker: bool) {
if self.is_released() {
self.request_close();
return;
}
let (start_tx, start_rx) = mpsc::channel();
let spawned = spawn_gated_close(
self.backend.clone(),
self.hpcon,
"conpty-oxide-input-close",
start_rx,
spawn_worker,
);
if let Err(err) = spawned {
log_close_worker_spawn_failure(&err, "conpty-oxide-input-close", true);
return;
}
let mut state = self.lock();
if state.close == CloseState::Done {
drop(state);
let _cancelled = start_tx.send(false);
return;
}
if state.released {
drop(state);
let _cancelled = start_tx.send(false);
self.request_close();
return;
}
state.close = CloseState::Done;
drop(state);
if start_tx.send(true).is_err() {
let mut state = self.lock();
debug_assert_eq!(state.close, CloseState::Done);
state.close = CloseState::NotRequested;
drop(state);
log_close_worker_spawn_failure(
&io::Error::new(
io::ErrorKind::BrokenPipe,
"pseudoconsole close worker exited before ownership handoff",
),
"conpty-oxide-input-close",
true,
);
}
}
fn close_if_due(&self, state: MutexGuard<'_, State>) {
if state.close == CloseState::Requested && state.reader != ReaderState::Open {
self.execute_close(state);
}
}
fn execute_close(&self, mut state: MutexGuard<'_, State>) {
debug_assert_ne!(state.close, CloseState::Done);
state.close = CloseState::Done;
drop(state);
unsafe { self.backend.close(self.hpcon) };
}
pub(crate) fn resize(&self, size: Size) -> io::Result<()> {
let mut state = self.lock();
if state.close == CloseState::Done {
return Err(io::Error::new(
io::ErrorKind::NotConnected,
"the pseudoconsole has been closed",
));
}
match unsafe { self.backend.resize(self.hpcon, size) } {
Ok(()) => {
state.size = size;
Ok(())
},
Err(err) => Err(normalize_session_end(err)),
}
}
#[cfg(any(feature = "blocking", feature = "tokio"))]
pub(crate) fn size(&self) -> Size {
self.lock().size
}
#[cfg(any(feature = "blocking", feature = "tokio"))]
pub(crate) fn supports_clear(&self) -> bool {
self.backend.supports_clear()
}
#[cfg(any(feature = "blocking", feature = "tokio"))]
#[cfg(test)]
pub(crate) fn supports_release(&self) -> bool {
self.backend.supports_release()
}
#[cfg(any(feature = "blocking", feature = "tokio"))]
pub(crate) fn backend_kind(&self) -> &BackendKind {
self.backend.kind()
}
pub(crate) fn clear(&self) -> io::Result<()> {
let state = self.lock();
if state.close == CloseState::Done {
return Err(io::Error::new(
io::ErrorKind::NotConnected,
"the pseudoconsole has been closed",
));
}
unsafe { self.backend.clear(self.hpcon) }.map_or_else(
|| {
Err(io::Error::new(
io::ErrorKind::Unsupported,
"this ConPTY backend does not export ClearPseudoConsole",
))
},
|result| result.map_err(normalize_session_end),
)
}
pub(crate) fn is_released(&self) -> bool {
self.lock().released
}
#[cfg(test)]
pub(crate) fn release_failed(&self) -> bool {
self.lock().release_failed
}
#[cfg(test)]
pub(crate) fn is_closed(&self) -> bool {
self.lock().close == CloseState::Done
}
pub(crate) fn reader_finished(&self) -> bool {
self.lock().reader != ReaderState::Open
}
}
impl Drop for ConsoleShared {
fn drop(&mut self) {
let state = self.state.get_mut().unwrap_or_else(PoisonError::into_inner);
if state.close == CloseState::Done {
return;
}
state.close = CloseState::Done;
if !state.released && state.reader != ReaderState::Drained {
#[cfg(test)]
if let Some(observer) = &state.drop_observer {
observer.store(1, Ordering::SeqCst);
}
spawn_detached_close(self.backend.clone(), self.hpcon, "conpty-oxide-close");
return;
}
#[cfg(test)]
if let Some(observer) = &state.drop_observer {
observer.store(2, Ordering::SeqCst);
}
unsafe { self.backend.close(self.hpcon) };
}
}
fn spawn_detached_close(backend: ConPtyBackend, hpcon: HPCON, name: &'static str) {
let spawned = thread::Builder::new().name(name.into()).spawn(move || {
unsafe { backend.close(hpcon) };
});
match spawned {
Ok(worker) => drop(worker),
Err(err) => log_close_worker_spawn_failure(&err, name, false),
}
}
#[cfg(any(feature = "blocking", feature = "tokio"))]
fn spawn_gated_close(
backend: ConPtyBackend,
hpcon: HPCON,
name: &'static str,
start: mpsc::Receiver<bool>,
spawn_worker: bool,
) -> io::Result<()> {
let spawned = if spawn_worker {
thread::Builder::new().name(name.into()).spawn(move || {
if matches!(start.recv(), Ok(true)) {
unsafe { backend.close(hpcon) };
}
})
} else {
Err(io::Error::other(
"pseudoconsole close worker spawn failure injected by a crate-local test",
))
};
spawned.map(drop)
}
#[cfg(feature = "tracing")]
fn log_close_worker_spawn_failure(err: &io::Error, worker: &'static str, retryable: bool) {
let logged = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
tracing::error!(
error = %err,
worker,
retryable,
"failed to spawn or start pseudoconsole close worker"
);
}));
if let Err(payload) = logged {
std::mem::forget(payload);
}
}
#[cfg(not(feature = "tracing"))]
const fn log_close_worker_spawn_failure(_err: &io::Error, _worker: &'static str, _retryable: bool) {
}
fn normalize_session_end(err: io::Error) -> io::Error {
if crate::core::is_disconnect_error(&err) {
io::Error::new(io::ErrorKind::NotConnected, err)
} else {
err
}
}
#[derive(Debug)]
pub(crate) struct PseudoConsole {
shared: Arc<ConsoleShared>,
}
impl PseudoConsole {
pub(crate) fn new(
backend: ConPtyBackend,
size: Size,
conin_read: OwnedHandle,
conout_write: OwnedHandle,
inherit_cursor: bool,
) -> io::Result<Self> {
let flags = if inherit_cursor {
PSEUDOCONSOLE_INHERIT_CURSOR
} else {
0
};
let hpcon = backend.create(
size,
conin_read.as_handle(),
conout_write.as_handle(),
flags,
)?;
drop(conin_read);
drop(conout_write);
Ok(Self {
shared: Arc::new(ConsoleShared {
backend,
hpcon,
state: Mutex::new(State::initial(size)),
}),
})
}
pub(crate) const fn shared(&self) -> &Arc<ConsoleShared> {
&self.shared
}
pub(super) fn spawn_capability(&self) -> io::Result<SpawnCapability<'_>> {
let state = self.shared.lock();
if state.close == CloseState::Done {
return Err(io::Error::new(
io::ErrorKind::NotConnected,
"the pseudoconsole has been closed",
));
}
Ok(SpawnCapability {
hpcon: self.shared.hpcon,
_state: state,
})
}
#[cfg(test)]
pub(crate) fn hpcon(&self) -> HPCON {
self.shared.hpcon
}
pub(crate) fn release_after_spawn(&self) -> io::Result<bool> {
self.shared.release_after_spawn()
}
pub(crate) fn resize(&self, size: Size) -> io::Result<()> {
self.shared.resize(size)
}
#[cfg(any(feature = "blocking", feature = "tokio"))]
pub(crate) fn size(&self) -> Size {
self.shared.size()
}
pub(crate) fn clear(&self) -> io::Result<()> {
self.shared.clear()
}
#[cfg(any(feature = "blocking", feature = "tokio"))]
pub(crate) fn supports_clear(&self) -> bool {
self.shared.supports_clear()
}
#[cfg(any(feature = "blocking", feature = "tokio"))]
#[cfg(test)]
pub(crate) fn supports_release(&self) -> bool {
self.shared.supports_release()
}
#[cfg(any(feature = "blocking", feature = "tokio"))]
pub(crate) fn backend_kind(&self) -> &BackendKind {
self.shared.backend_kind()
}
#[cfg(test)]
pub(crate) fn is_released(&self) -> bool {
self.shared.is_released()
}
}
#[cfg(test)]
#[path = "pseudocon_tests.rs"]
mod tests;