#![allow(dead_code)]
use std::collections::VecDeque;
use std::fmt;
use std::fs::File;
use std::io::{self, Read, Write};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::Duration;
use tokio::sync::Notify;
use super::errno::{code_of_raw_os_error, errno_message, errno_of};
pub(crate) mod redirect_flags {
pub(crate) const STDIN: u8 = 1;
pub(crate) const STDOUT: u8 = 2;
pub(crate) const STDERR: u8 = 4;
pub(crate) const APPEND: u8 = 8;
pub(crate) const DUPLICATE_OUT: u8 = 16;
}
pub(crate) const HIGH_WATER: usize = 64 * 1024;
pub(crate) const READ_CHUNK: usize = 64 * 1024;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum Which {
Stdout,
Stderr,
}
impl Which {
pub(crate) fn as_str(self) -> &'static str {
match self {
Which::Stdout => "stdout",
Which::Stderr => "stderr",
}
}
}
pub(crate) fn redirects_elsewhere(flags: u8, which: Which) -> bool {
let bit = match which {
Which::Stdout => redirect_flags::STDOUT,
Which::Stderr => redirect_flags::STDERR,
};
if flags & redirect_flags::DUPLICATE_OUT != 0 {
flags & bit == 0
} else {
flags & bit != 0
}
}
#[derive(Clone, Debug, Default)]
pub(crate) struct SharedBuf(Arc<Mutex<Vec<u8>>>);
impl SharedBuf {
pub(crate) fn new() -> Self {
Self::default()
}
fn lock(&self) -> MutexGuard<'_, Vec<u8>> {
self.0.lock().unwrap_or_else(|e| e.into_inner())
}
pub(crate) fn append(&self, bytes: &[u8]) {
if !bytes.is_empty() {
self.lock().extend_from_slice(bytes);
}
}
pub(crate) fn len(&self) -> usize {
self.lock().len()
}
pub(crate) fn is_empty(&self) -> bool {
self.lock().is_empty()
}
pub(crate) fn to_vec(&self) -> Vec<u8> {
self.lock().clone()
}
pub(crate) fn take(&self) -> Vec<u8> {
std::mem::take(&mut *self.lock())
}
pub(crate) fn clear(&self) {
self.lock().clear();
}
pub(crate) fn ptr_eq(&self, other: &Self) -> bool {
Arc::ptr_eq(&self.0, &other.0)
}
pub(crate) fn with<R>(&self, f: impl FnOnce(&mut Vec<u8>) -> R) -> R {
f(&mut self.lock())
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct ShellSysError {
pub(crate) code: &'static str,
pub(crate) errno: i32,
pub(crate) path: String,
pub(crate) syscall: &'static str,
}
impl ShellSysError {
pub(crate) fn new(code: &'static str) -> Self {
Self {
code,
errno: errno_of(code),
path: String::new(),
syscall: "",
}
}
pub(crate) fn with_path(mut self, path: impl Into<String>) -> Self {
self.path = path.into();
self
}
pub(crate) fn with_syscall(mut self, syscall: &'static str) -> Self {
self.syscall = syscall;
self
}
pub(crate) fn epipe() -> Self {
Self::new("EPIPE")
}
pub(crate) fn message(&self) -> &'static str {
errno_message(self.code).unwrap_or(self.code)
}
pub(crate) fn display(&self) -> String {
format!("bun: {}: {}", self.message(), self.path)
}
pub(crate) fn from_io(err: &io::Error, path: impl Into<String>) -> Self {
let (code, errno) = match err.raw_os_error() {
Some(raw) => {
let code = code_of_raw_os_error(raw).unwrap_or("UNKNOWN");
let errno = if cfg!(windows) {
match errno_of(code) {
0 => raw.abs(),
n => n,
}
} else {
raw.abs()
};
(code, errno)
}
None => {
let code = match err.kind() {
io::ErrorKind::NotFound => "ENOENT",
io::ErrorKind::PermissionDenied => "EACCES",
io::ErrorKind::AlreadyExists => "EEXIST",
io::ErrorKind::BrokenPipe => "EPIPE",
io::ErrorKind::WouldBlock => "EAGAIN",
io::ErrorKind::InvalidInput => "EINVAL",
io::ErrorKind::NotADirectory => "ENOTDIR",
io::ErrorKind::IsADirectory => "EISDIR",
io::ErrorKind::DirectoryNotEmpty => "ENOTEMPTY",
io::ErrorKind::StorageFull => "ENOSPC",
_ => "EIO",
};
(code, errno_of(code))
}
};
Self {
code,
errno,
path: path.into(),
syscall: "",
}
}
}
impl fmt::Display for ShellSysError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.message())
}
}
impl std::error::Error for ShellSysError {}
#[derive(Default)]
struct ChannelState {
chunks: VecDeque<Vec<u8>>,
pending: usize,
write_closed: bool,
read_closed: bool,
drain_gen: u64,
}
#[derive(Default)]
struct ChannelInner {
state: Mutex<ChannelState>,
readable: Notify,
drained: Notify,
}
#[derive(Clone, Default)]
pub(crate) struct Channel(Arc<ChannelInner>);
impl fmt::Debug for Channel {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let st = self.state();
f.debug_struct("Channel")
.field("pending", &st.pending)
.field("write_closed", &st.write_closed)
.field("read_closed", &st.read_closed)
.finish()
}
}
impl Channel {
pub(crate) fn new() -> Self {
Self::default()
}
fn state(&self) -> MutexGuard<'_, ChannelState> {
self.0.state.lock().unwrap_or_else(|e| e.into_inner())
}
pub(crate) async fn write(&self, bytes: Vec<u8>) -> Result<(), ShellSysError> {
let gen = {
let mut st = self.state();
if st.read_closed {
return Err(ShellSysError::epipe());
}
st.pending += bytes.len();
st.chunks.push_back(bytes);
self.0.readable.notify_waiters();
if st.pending <= HIGH_WATER {
return Ok(());
}
st.drain_gen
};
loop {
let notified = self.0.drained.notified();
tokio::pin!(notified);
notified.as_mut().enable();
{
let st = self.state();
if st.drain_gen != gen {
return Ok(());
}
if st.read_closed {
return Err(ShellSysError::epipe());
}
}
notified.await;
}
}
pub(crate) fn close(&self) {
self.state().write_closed = true;
self.0.readable.notify_waiters();
}
pub(crate) fn close_read(&self) {
{
let mut st = self.state();
if st.read_closed {
return;
}
st.read_closed = true;
st.chunks.clear();
st.pending = 0;
}
self.0.drained.notify_waiters();
self.0.readable.notify_waiters();
}
pub(crate) fn is_read_closed(&self) -> bool {
self.state().read_closed
}
pub(crate) async fn read(&self) -> Option<Vec<u8>> {
loop {
let notified = self.0.readable.notified();
tokio::pin!(notified);
notified.as_mut().enable();
{
let mut st = self.state();
if let Some(chunk) = st.chunks.pop_front() {
st.pending -= chunk.len();
if st.pending <= HIGH_WATER {
st.drain_gen = st.drain_gen.wrapping_add(1);
self.0.drained.notify_waiters();
}
return Some(chunk);
}
if st.write_closed || st.read_closed {
return None;
}
}
notified.await;
}
}
}
pub(crate) enum WriterTarget {
File(Arc<File>),
Stdout,
Stderr,
Channel(Channel),
}
struct WriterInner {
target: WriterTarget,
err: tokio::sync::Mutex<Option<ShellSysError>>,
closed: AtomicBool,
}
impl Drop for WriterInner {
fn drop(&mut self) {
if let WriterTarget::Channel(c) = &self.target {
c.close();
}
}
}
#[derive(Clone)]
pub(crate) struct Writer(Arc<WriterInner>);
impl fmt::Debug for Writer {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let kind = match &self.0.target {
WriterTarget::File(_) => "file",
WriterTarget::Stdout => "stdout",
WriterTarget::Stderr => "stderr",
WriterTarget::Channel(_) => "channel",
};
write!(f, "Writer({kind})")
}
}
impl Writer {
pub(crate) fn new(target: WriterTarget) -> Self {
Self(Arc::new(WriterInner {
target,
err: tokio::sync::Mutex::new(None),
closed: AtomicBool::new(false),
}))
}
pub(crate) fn file(file: File) -> Self {
Self::new(WriterTarget::File(Arc::new(file)))
}
pub(crate) fn stdout() -> Self {
Self::new(WriterTarget::Stdout)
}
pub(crate) fn stderr() -> Self {
Self::new(WriterTarget::Stderr)
}
pub(crate) fn channel(channel: Channel) -> Self {
Self::new(WriterTarget::Channel(channel))
}
pub(crate) fn target(&self) -> &WriterTarget {
&self.0.target
}
pub(crate) fn is_channel(&self) -> bool {
matches!(self.0.target, WriterTarget::Channel(_))
}
pub(crate) fn ptr_eq(&self, other: &Self) -> bool {
Arc::ptr_eq(&self.0, &other.0)
}
pub(crate) async fn write(
&self,
data: &[u8],
captured: Option<&SharedBuf>,
) -> Result<(), ShellSysError> {
let mut err = self.0.err.lock().await;
if let Some(e) = &*err {
return Err(e.clone());
}
if data.is_empty() {
return Ok(());
}
let res = match &self.0.target {
WriterTarget::File(f) => {
let f = Arc::clone(f);
let data = data.to_vec();
run_blocking(move || write_all_retrying(&mut &*f, &data)).await
}
WriterTarget::Stdout => {
let data = data.to_vec();
run_blocking(move || {
let mut out = io::stdout().lock();
write_all_retrying(&mut out, &data)?;
out.flush()
})
.await
}
WriterTarget::Stderr => {
let data = data.to_vec();
run_blocking(move || {
let mut out = io::stderr().lock();
write_all_retrying(&mut out, &data)?;
out.flush()
})
.await
}
WriterTarget::Channel(c) => c
.write(data.to_vec())
.await
.map_err(|_| io::Error::from(io::ErrorKind::BrokenPipe)),
};
match res {
Ok(()) => {
if let Some(c) = captured {
c.append(data);
}
Ok(())
}
Err(e) => {
let e = ShellSysError::from_io(&e, "").with_syscall("write");
*err = Some(e.clone());
Err(e)
}
}
}
pub(crate) async fn close(&self) {
let _guard = self.0.err.lock().await;
if !self.0.closed.swap(true, Ordering::SeqCst) {
if let WriterTarget::Channel(c) = &self.0.target {
c.close();
}
}
}
pub(crate) fn try_clone_file(&self) -> Option<io::Result<File>> {
match &self.0.target {
WriterTarget::File(f) => Some(f.try_clone()),
WriterTarget::Stdout => Some(dup_process_stdio(Which::Stdout)),
WriterTarget::Stderr => Some(dup_process_stdio(Which::Stderr)),
WriterTarget::Channel(_) => None,
}
}
}
fn write_all_retrying(w: &mut impl Write, mut data: &[u8]) -> io::Result<()> {
while !data.is_empty() {
match w.write(data) {
Ok(0) => return Err(io::Error::from(io::ErrorKind::WriteZero)),
Ok(n) => data = &data[n..],
Err(e) if e.kind() == io::ErrorKind::Interrupted => {}
Err(e) if e.kind() == io::ErrorKind::WouldBlock => {
std::thread::sleep(Duration::from_millis(1));
}
Err(e) => return Err(e),
}
}
Ok(())
}
async fn run_blocking<T: Send + 'static>(
f: impl FnOnce() -> io::Result<T> + Send + 'static,
) -> io::Result<T> {
match tokio::task::spawn_blocking(f).await {
Ok(r) => r,
Err(e) => Err(io::Error::other(e)),
}
}
pub(crate) fn dup_process_stdio(which: Which) -> io::Result<File> {
dup_std(Some(which))
}
fn dup_process_stdin() -> io::Result<File> {
dup_std(None)
}
#[cfg(unix)]
fn dup_std(which: Option<Which>) -> io::Result<File> {
use std::os::fd::AsFd;
let owned = match which {
None => io::stdin().as_fd().try_clone_to_owned()?,
Some(Which::Stdout) => io::stdout().as_fd().try_clone_to_owned()?,
Some(Which::Stderr) => io::stderr().as_fd().try_clone_to_owned()?,
};
Ok(File::from(owned))
}
#[cfg(windows)]
fn dup_std(which: Option<Which>) -> io::Result<File> {
use std::os::windows::io::AsHandle;
let owned = match which {
None => io::stdin().as_handle().try_clone_to_owned()?,
Some(Which::Stdout) => io::stdout().as_handle().try_clone_to_owned()?,
Some(Which::Stderr) => io::stderr().as_handle().try_clone_to_owned()?,
};
Ok(File::from(owned))
}
#[cfg(not(any(unix, windows)))]
fn dup_std(_which: Option<Which>) -> io::Result<File> {
Err(io::Error::from(io::ErrorKind::Unsupported))
}
pub(crate) enum ReaderSource {
Channel(Channel),
File(Arc<File>),
Stdin,
}
struct ReaderInner {
source: ReaderSource,
stdin: Mutex<Option<Arc<File>>>,
}
impl Drop for ReaderInner {
fn drop(&mut self) {
if let ReaderSource::Channel(c) = &self.source {
c.close_read();
}
}
}
#[derive(Clone)]
pub(crate) struct Reader(Arc<ReaderInner>);
impl fmt::Debug for Reader {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let kind = match &self.0.source {
ReaderSource::Channel(_) => "channel",
ReaderSource::File(_) => "file",
ReaderSource::Stdin => "stdin",
};
write!(f, "Reader({kind})")
}
}
impl Reader {
pub(crate) fn new(source: ReaderSource) -> Self {
Self(Arc::new(ReaderInner {
source,
stdin: Mutex::new(None),
}))
}
pub(crate) fn channel(channel: Channel) -> Self {
Self::new(ReaderSource::Channel(channel))
}
pub(crate) fn file(file: File) -> Self {
Self::new(ReaderSource::File(Arc::new(file)))
}
pub(crate) fn stdin() -> Self {
Self::new(ReaderSource::Stdin)
}
pub(crate) fn source(&self) -> &ReaderSource {
&self.0.source
}
pub(crate) fn ptr_eq(&self, other: &Self) -> bool {
Arc::ptr_eq(&self.0, &other.0)
}
pub(crate) async fn read_chunk(&self) -> Result<Option<Vec<u8>>, ShellSysError> {
let file = match &self.0.source {
ReaderSource::Channel(c) => return Ok(c.read().await),
ReaderSource::File(f) => Arc::clone(f),
ReaderSource::Stdin => self.stdin_file()?,
};
let detached = matches!(self.0.source, ReaderSource::Stdin);
loop {
match read_once(Arc::clone(&file), detached).await {
Ok(buf) if buf.is_empty() => return Ok(None),
Ok(buf) => return Ok(Some(buf)),
Err(e) if e.kind() == io::ErrorKind::Interrupted => {}
Err(e) if e.kind() == io::ErrorKind::WouldBlock => {
tokio::time::sleep(Duration::from_millis(5)).await;
}
Err(e) if e.kind() == io::ErrorKind::BrokenPipe => return Ok(None),
Err(e) => return Err(ShellSysError::from_io(&e, "").with_syscall("read")),
}
}
}
pub(crate) async fn read_to_end(&self) -> Result<Vec<u8>, ShellSysError> {
let mut out = Vec::new();
while let Some(chunk) = self.read_chunk().await? {
out.extend_from_slice(&chunk);
}
Ok(out)
}
pub(crate) fn close(&self) {
if let ReaderSource::Channel(c) = &self.0.source {
c.close_read();
}
}
pub(crate) fn try_clone_file(&self) -> Option<io::Result<File>> {
match &self.0.source {
ReaderSource::File(f) => Some(f.try_clone()),
_ => None,
}
}
fn stdin_file(&self) -> Result<Arc<File>, ShellSysError> {
let mut guard = self.0.stdin.lock().unwrap_or_else(|e| e.into_inner());
if let Some(f) = &*guard {
return Ok(Arc::clone(f));
}
let f = Arc::new(
dup_process_stdin().map_err(|e| ShellSysError::from_io(&e, "").with_syscall("read"))?,
);
*guard = Some(Arc::clone(&f));
Ok(f)
}
}
async fn read_once(file: Arc<File>, detached: bool) -> io::Result<Vec<u8>> {
let read = move || {
let mut buf = vec![0u8; READ_CHUNK];
let n = (&*file).read(&mut buf)?;
buf.truncate(n);
Ok(buf)
};
if !detached {
return run_blocking(read).await;
}
let (tx, rx) = tokio::sync::oneshot::channel();
std::thread::Builder::new()
.name("bun-shell-stdin".into())
.spawn(move || {
let _ = tx.send(read());
})?;
rx.await
.unwrap_or_else(|_| Err(io::Error::from(io::ErrorKind::BrokenPipe)))
}
#[derive(Clone, Debug)]
pub(crate) enum OutKind {
Fd {
writer: Writer,
captured: Option<SharedBuf>,
},
Pipe,
Ignore,
}
impl OutKind {
pub(crate) fn fd(writer: Writer, captured: Option<SharedBuf>) -> Self {
OutKind::Fd { writer, captured }
}
pub(crate) async fn write(
&self,
data: &[u8],
buffered: &SharedBuf,
) -> Result<(), ShellSysError> {
match self {
OutKind::Fd { writer, captured } => writer.write(data, captured.as_ref()).await,
OutKind::Pipe => {
buffered.append(data);
Ok(())
}
OutKind::Ignore => Ok(()),
}
}
}
#[derive(Clone, Debug)]
pub(crate) enum InKind {
Fd(Reader),
Ignore,
}
#[derive(Clone, Debug)]
pub(crate) struct ShellIO {
pub(crate) stdin: InKind,
pub(crate) stdout: OutKind,
pub(crate) stderr: OutKind,
}
impl ShellIO {
pub(crate) fn out(&self, which: Which) -> &OutKind {
match which {
Which::Stdout => &self.stdout,
Which::Stderr => &self.stderr,
}
}
}
#[cfg(test)]
mod tests;