use std::sync::{
Arc, Mutex,
atomic::{AtomicBool, Ordering},
};
#[cfg_attr(miri, allow(unused_imports))]
use anyhow::Result;
use crate::backend::{PipeInner, ScriptPipe};
pub type SharedInput = Arc<Mutex<dyn std::io::Read + Send>>;
pub type SharedOutput = Arc<Mutex<dyn std::io::Write + Send>>;
#[cfg(not(miri))]
#[derive(Clone)]
pub struct OsPipeReader {
inner: Arc<Mutex<Option<std::io::PipeReader>>>,
}
#[cfg(not(miri))]
#[derive(Clone)]
pub struct OsPipeWriter {
inner: Arc<Mutex<Option<std::io::PipeWriter>>>,
}
#[cfg(not(miri))]
impl OsPipeReader {
fn new(reader: std::io::PipeReader) -> Self {
Self {
inner: Arc::new(Mutex::new(Some(reader))),
}
}
pub fn take(&self) -> Result<std::io::PipeReader> {
self.inner
.lock()
.map_err(|_| anyhow::anyhow!("os pipe reader lock poisoned"))?
.take()
.ok_or_else(|| {
anyhow::anyhow!("os pipe handle has already been consumed by another process")
})
}
pub fn is_consumed(&self) -> bool {
self.inner.lock().map(|g| g.is_none()).unwrap_or(false)
}
}
#[cfg(not(miri))]
impl OsPipeWriter {
fn new(writer: std::io::PipeWriter) -> Self {
Self {
inner: Arc::new(Mutex::new(Some(writer))),
}
}
pub fn take(&self) -> Result<std::io::PipeWriter> {
self.inner
.lock()
.map_err(|_| anyhow::anyhow!("os pipe writer lock poisoned"))?
.take()
.ok_or_else(|| {
anyhow::anyhow!("os pipe handle has already been consumed by another process")
})
}
pub fn is_consumed(&self) -> bool {
self.inner.lock().map(|g| g.is_none()).unwrap_or(false)
}
}
#[cfg(not(miri))]
pub fn create_os_pipe() -> Result<(OsPipeReader, OsPipeWriter)> {
let (reader, writer) = std::io::pipe()?;
Ok((OsPipeReader::new(reader), OsPipeWriter::new(writer)))
}
#[cfg(not(miri))]
#[derive(Clone)]
pub struct OsPipeEntry {
pub writer: OsPipeWriter,
pub reader: OsPipeReader,
}
#[cfg(not(miri))]
impl OsPipeEntry {
pub fn new() -> anyhow::Result<Self> {
let (reader, writer) = create_os_pipe()?;
Ok(Self { writer, reader })
}
pub fn is_spent(&self) -> bool {
self.reader.is_consumed() && self.writer.is_consumed()
}
}
pub enum Slot {
Unbound,
Script {
backend: Arc<PipeInner>,
},
#[cfg(not(miri))]
Os {
entry: OsPipeEntry,
},
}
impl std::fmt::Debug for Slot {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Slot::Unbound => write!(f, "Unbound"),
Slot::Script { .. } => write!(f, "Script(..)"),
#[cfg(not(miri))]
Slot::Os { .. } => write!(f, "Os(..)"),
}
}
}
#[derive(Clone, Debug)]
pub struct PipeHandle {
cell: Arc<Mutex<Slot>>,
declaring_task_id: u64,
escaped_to_child: Arc<AtomicBool>,
}
pub fn new_handle_in_task(task_id: u64) -> PipeHandle {
PipeHandle {
cell: Arc::new(Mutex::new(Slot::Unbound)),
declaring_task_id: task_id,
escaped_to_child: Arc::new(AtomicBool::new(false)),
}
}
impl PipeHandle {
pub fn declaring_task(&self) -> u64 {
self.declaring_task_id
}
pub fn has_escaped(&self) -> bool {
self.escaped_to_child.load(Ordering::SeqCst)
}
pub fn mark_escaped(&self) {
self.escaped_to_child.store(true, Ordering::SeqCst);
}
pub(crate) fn cell(&self) -> &Arc<Mutex<Slot>> {
&self.cell
}
pub fn ptr_eq(&self, other: &Self) -> bool {
Arc::ptr_eq(&self.cell, &other.cell)
}
}
pub enum Materialized {
Script(Arc<PipeInner>),
#[cfg(not(miri))]
Os(OsPipeEntry),
}
pub fn materialize(handle: &PipeHandle, promote: bool) -> anyhow::Result<Materialized> {
let mut guard = handle
.cell()
.lock()
.map_err(|_| anyhow::anyhow!("pipe handle lock poisoned"))?;
match &*guard {
Slot::Script { backend } => Ok(Materialized::Script(Arc::clone(backend))),
#[cfg(not(miri))]
Slot::Os { entry } => Ok(Materialized::Os(entry.clone())),
Slot::Unbound => {
#[cfg(not(miri))]
if promote {
let entry = OsPipeEntry::new()?;
let out = entry.clone();
*guard = Slot::Os { entry };
return Ok(Materialized::Os(out));
}
#[cfg(miri)]
let _ = promote;
let pipe = ScriptPipe::new();
let backend = pipe.pipe_inner();
*guard = Slot::Script {
backend: Arc::clone(&backend),
};
Ok(Materialized::Script(backend))
}
}
}