use std::io::{self, Read, Write};
use std::sync::{Arc, Condvar, Mutex};
use std::time::Duration;
use crate::slot::{SharedInput, SharedOutput};
use crate::spill::SpillBuffer;
pub struct ScriptPipe {
inner: Arc<PipeInner>,
reader: SharedInput,
}
impl ScriptPipe {
pub fn new() -> Self {
let inner = Arc::new(PipeInner::new());
let reader: SharedInput = Arc::new(Mutex::new(PipeReader::new(inner.clone())));
Self { inner, reader }
}
pub fn with_thresholds(spill_threshold: usize, max_backlog: u64) -> Self {
let inner = Arc::new(PipeInner::with_thresholds(spill_threshold, max_backlog));
let reader: SharedInput = Arc::new(Mutex::new(PipeReader::new(inner.clone())));
Self { inner, reader }
}
pub fn reader(&self) -> SharedInput {
self.reader.clone()
}
pub fn endpoint(&self) -> ScriptPipeEndpoint {
ScriptPipeEndpoint::new(self.inner.clone())
}
pub fn pipe_inner(&self) -> Arc<PipeInner> {
self.inner.clone()
}
#[allow(clippy::disallowed_types)]
pub fn temp_path(&self) -> Option<std::path::PathBuf> {
self.inner.temp_path()
}
}
impl Default for ScriptPipe {
fn default() -> Self {
Self::new()
}
}
#[derive(Clone)]
pub struct ScriptPipeEndpoint {
inner: Arc<PipeInner>,
}
impl ScriptPipeEndpoint {
fn new(inner: Arc<PipeInner>) -> Self {
Self { inner }
}
pub fn for_backend(inner: &Arc<PipeInner>) -> Self {
Self::new(Arc::clone(inner))
}
pub fn stream_handle(&self) -> SharedOutput {
Arc::new(Mutex::new(PipeWriter::new(self.inner.clone())))
}
}
pub struct PipeInner {
state: Mutex<PipeState>,
ready: Condvar,
}
struct PipeState {
buffer: SpillBuffer,
writers: usize,
keepers: usize,
closed: bool,
}
impl PipeState {
fn new(spill_threshold: usize, max_backlog: u64) -> Self {
Self {
buffer: SpillBuffer::with_thresholds(spill_threshold, max_backlog),
writers: 0,
keepers: 0,
closed: false,
}
}
}
impl PipeInner {
fn new() -> Self {
Self {
state: Mutex::new(PipeState::new(
crate::spill::DEFAULT_SPILL_THRESHOLD,
crate::spill::DEFAULT_MAX_BACKLOG,
)),
ready: Condvar::new(),
}
}
pub fn reader_handle(self: &Arc<Self>) -> crate::slot::SharedInput {
Arc::new(Mutex::new(PipeReader::new(Arc::clone(self))))
}
pub fn writer_handle(self: &Arc<Self>) -> crate::slot::SharedOutput {
Arc::new(Mutex::new(PipeWriter::new(Arc::clone(self))))
}
fn with_thresholds(spill_threshold: usize, max_backlog: u64) -> Self {
Self {
state: Mutex::new(PipeState::new(spill_threshold, max_backlog)),
ready: Condvar::new(),
}
}
#[allow(clippy::disallowed_types)]
fn temp_path(&self) -> Option<std::path::PathBuf> {
self.lock_state().buffer.temp_path()
}
fn attach_writer(&self) {
let mut state = self.lock_state();
state.writers += 1;
state.closed = false;
}
pub fn writer_count(&self) -> usize {
self.lock_state().writers
}
pub fn buffered_bytes(&self) -> u64 {
self.lock_state().buffer.buffered_bytes()
}
pub fn peek_bytes(&self) -> io::Result<Vec<u8>> {
self.lock_state().buffer.peek_bytes()
}
fn detach_writer(&self) {
let mut state = self.lock_state();
state.writers = state.writers.saturating_sub(1);
if state.writers == 0 && state.keepers == 0 {
state.closed = true;
}
drop(state);
self.ready.notify_all();
}
pub fn force_close(&self) {
let mut state = self.lock_state();
state.closed = true;
drop(state);
self.ready.notify_all();
}
pub fn pin_keeper(&self) {
let mut state = self.lock_state();
state.keepers += 1;
}
pub fn unpin_keeper(&self) {
let mut state = self.lock_state();
state.keepers = state.keepers.saturating_sub(1);
if state.writers == 0 && state.keepers == 0 {
state.closed = true;
}
drop(state);
self.ready.notify_all();
}
fn push_bytes(&self, data: &[u8]) -> io::Result<()> {
let state = self.lock_state();
let res = state.buffer.push_bytes(data);
drop(state);
self.ready.notify_all();
res
}
fn read_into(&self, buf: &mut [u8]) -> io::Result<usize> {
if buf.is_empty() {
return Ok(0);
}
let mut state = self.lock_state();
loop {
let n = state.buffer.read_into(buf)?;
if n > 0 {
return Ok(n);
}
if state.closed {
return Ok(0);
}
state = self
.ready
.wait(state)
.map_err(|_| io::Error::other("pipe wait poisoned"))?;
}
}
pub fn read_into_timeout(
&self,
buf: &mut [u8],
backstop: Duration,
) -> io::Result<Option<usize>> {
if buf.is_empty() {
return Ok(Some(0));
}
let mut state = self.lock_state();
loop {
let n = state.buffer.read_into(buf)?;
if n > 0 {
return Ok(Some(n));
}
if state.closed {
return Ok(Some(0));
}
let (guard, waited) = self
.ready
.wait_timeout(state, backstop)
.map_err(|_| io::Error::other("pipe wait poisoned"))?;
state = guard;
if waited.timed_out() {
return Ok(None);
}
}
}
fn lock_state(&self) -> std::sync::MutexGuard<'_, PipeState> {
self.state.lock().expect("script pipe state poisoned")
}
}
struct PipeReader {
inner: Arc<PipeInner>,
}
impl PipeReader {
fn new(inner: Arc<PipeInner>) -> Self {
Self { inner }
}
}
impl Read for PipeReader {
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
self.inner.read_into(buf)
}
}
struct PipeWriter {
inner: Arc<PipeInner>,
}
impl PipeWriter {
fn new(inner: Arc<PipeInner>) -> Self {
inner.attach_writer();
Self { inner }
}
}
impl Write for PipeWriter {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
self.inner.push_bytes(buf)?;
Ok(buf.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
impl Drop for PipeWriter {
fn drop(&mut self) {
self.inner.detach_writer();
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PipeKindDesc {
Unbound,
Script,
#[cfg_attr(miri, allow(dead_code))]
Os,
}
impl PipeKindDesc {
pub fn as_str(self) -> &'static str {
match self {
PipeKindDesc::Unbound => "unbound",
PipeKindDesc::Script => "script",
PipeKindDesc::Os => "os",
}
}
pub fn is_os(self) -> bool {
matches!(self, PipeKindDesc::Os)
}
}
#[derive(Debug, Clone)]
pub struct PipeInfo {
pub kind: PipeKindDesc,
pub buffered: u64,
pub readers: usize,
pub writers: usize,
}
pub fn script_backend(handle: &crate::slot::PipeHandle) -> Option<Arc<PipeInner>> {
use crate::slot::Slot;
let guard = handle.cell().lock().expect("pipe handle lock poisoned");
match &*guard {
Slot::Script { backend } => Some(Arc::clone(backend)),
Slot::Unbound => None,
#[cfg(not(miri))]
Slot::Os { .. } => None,
}
}
pub fn peek(handle: &crate::slot::PipeHandle) -> anyhow::Result<Vec<u8>> {
use crate::slot::Slot;
let backend = {
let guard = handle.cell().lock().expect("pipe handle lock poisoned");
match &*guard {
Slot::Script { backend } => Some(Arc::clone(backend)),
Slot::Unbound => {
return Err(anyhow::anyhow!(
"cannot peek unbound pipe: bind it to a command first"
));
}
#[cfg(not(miri))]
Slot::Os { .. } => {
return Err(anyhow::anyhow!(
"cannot peek OS-materialized pipe: drain it through a bound command instead"
));
}
}
};
let Some(inner) = backend else {
unreachable!("unbound/os arms return above");
};
inner
.peek_bytes()
.map_err(|e| anyhow::anyhow!("failed to peek pipe: {e}"))
}
pub fn inspect(handle: &crate::slot::PipeHandle) -> PipeInfo {
use crate::slot::Slot;
let backend = {
let guard = handle.cell().lock().expect("pipe handle lock poisoned");
match &*guard {
Slot::Unbound => None,
Slot::Script { backend } => Some(backend.clone()),
#[cfg(not(miri))]
Slot::Os { .. } => {
return PipeInfo {
kind: PipeKindDesc::Os,
buffered: 0,
readers: 1,
writers: 1,
};
}
}
};
match backend {
Some(inner) => PipeInfo {
kind: PipeKindDesc::Script,
buffered: inner.buffered_bytes(),
readers: 1,
writers: inner.writer_count(),
},
None => PipeInfo {
kind: PipeKindDesc::Unbound,
buffered: 0,
readers: 0,
writers: 0,
},
}
}
pub struct KeeperGuard {
inner: Option<Arc<PipeInner>>,
}
impl KeeperGuard {
pub fn new(inner: Arc<PipeInner>) -> Self {
inner.pin_keeper();
Self { inner: Some(inner) }
}
}
impl Drop for KeeperGuard {
fn drop(&mut self) {
if let Some(inner) = self.inner.take() {
inner.unpin_keeper();
}
}
}