Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
117 changes: 109 additions & 8 deletions src/uu/cat/src/cat.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ use std::os::fd::AsFd;
use std::os::unix::fs::FileTypeExt;
use thiserror::Error;
use uucore::display::Quotable;
use uucore::error::{UResult, strip_errno};
use uucore::error::{FromIo, UIoError, UResult, strip_errno};
use uucore::translate;
use uucore::{fast_inc::fast_inc_one, format_usage};

Expand Down Expand Up @@ -80,6 +80,10 @@ enum CatError {
/// Wrapper around `io::Error`
#[error("{}", strip_errno(.0))]
Io(#[from] io::Error),
/// A write error to the output; unlike [`Self::Io`] it is reported
/// without the input filename.
#[error("{0}")]
Write(Box<UIoError>),
Comment thread
MuntasirSZN marked this conversation as resolved.
/// Unknown file type; it's not a regular file, socket, etc.
#[error("{}", translate!("cat-error-unknown-filetype", "ft_debug" => .ft_debug))]
UnknownFiletype {
Expand All @@ -99,6 +103,25 @@ enum CatError {

type CatResult<T> = Result<T, CatError>;

impl CatError {
/// Compose the diagnostic for the input at `path`.
///
/// Write errors are reported without the filename,
/// while input errors name the file.
fn display_message(&self, path: &OsString) -> String {
Comment thread
MuntasirSZN marked this conversation as resolved.
if matches!(self, Self::Write(_)) {
format!("{self}")
} else {
format!("{}: {self}", path.maybe_quote())
}
}
}

/// An error writing to the output, using the shared uucore "write error"
/// context (as opposed to [`CatError::Io`], which names the input file).
fn write_err(err: io::Error) -> CatError {
CatError::Write(err.map_err_context(|| translate!("common-write-error")))
}
#[derive(PartialEq)]
enum NumberingMode {
None,
Expand Down Expand Up @@ -403,7 +426,7 @@ where

for path in files {
if let Err(err) = cat_path(path, options, &mut state) {
error_messages.push(format!("{}: {err}", path.maybe_quote()));
error_messages.push(err.display_message(path));
}
}
if state.skipped_carriage_return {
Expand Down Expand Up @@ -470,19 +493,96 @@ fn get_input_type(path: &OsString) -> CatResult<InputType> {
/// simple memory copy.
fn print_fast<R: FdReadable>(handle: &mut InputHandle<R>) -> CatResult<()> {
let stdout = io::stdout();
// Try to use the splice() system call for faster writing.
#[cfg(any(target_os = "linux", target_os = "android"))]
let mut stdout = stdout;
// Try to use the splice() system call for faster writing. If it works, we're done.
#[cfg(any(target_os = "linux", target_os = "android"))]
if uucore::pipes::splice_unbounded_auto(&handle.reader, &mut stdout)?.is_ok() {
return Ok(());
match splice_cat(&handle.reader, &stdout) {
// splice() copied the whole input.
Ok(()) => return Ok(()),
// splice() is unusable here; fall back on slower writing.
Err(CatError::Io(e)) if uucore::pipes::splice_unusable(&e) => {}
Err(e) => return Err(e),
}

// If we're not on Linux or Android, or the splice() call failed,
// fall back on slower writing.
print_unbuffered(handle, stdout)
}

/// Copy `source` to `stdout` with the `splice(2)` syscall, diagnosing
/// errors like GNU.
///
/// Returns `Ok(())` if splice handled the whole input. `Err(EINVAL)`
/// (checked with [`uucore::pipes::splice_unusable`]) means splice is
/// unusable here and the caller should fall back on read/write; ENOSYS is
/// folded into EINVAL, and so are failures before any byte was moved. Once
/// splice is known to work, every error is fatal, reported either as an
/// input error (naming the file) or as a write error.
#[cfg(any(target_os = "linux", target_os = "android"))]
fn splice_cat<R: FdReadable>(reader: &R, stdout: &io::Stdout) -> CatResult<()> {
use uucore::pipes::{MAX_ROOTLESS_PIPE_SIZE, pipe, splice};

// Create a broker pipe, since splice(2) needs a pipe on one side and we
// want to distinguish read errors from write errors. If the pipe cannot
// be created (e.g. file descriptor exhaustion), fall back on read/write.
let Ok((pipe_rd, pipe_wr)) = pipe::<false>() else {
return Err(CatError::Io(unusable_errno()));
};

// First input splice: any failure just selects read/write.
let mut remaining = match splice(reader, &pipe_wr, MAX_ROOTLESS_PIPE_SIZE) {
// End of input: splice handled the whole file.
Ok(0) => return Ok(()),
Ok(bytes_read) => bytes_read,
Err(_) => return Err(CatError::Io(unusable_errno())),
};

// First output splice: if stdout cannot accept splice data, drain the
// broker pipe with plain read/write, then let the caller continue
// likewise.
while remaining > 0 {
match splice(&pipe_rd, stdout, remaining) {
// No progress; stop splicing.
Ok(0) => return Err(CatError::Io(unusable_errno())),
Ok(bytes_written) => remaining -= bytes_written,
Err(_) => {
let mut drain = Vec::with_capacity(remaining);
let _ = pipe_rd.take(remaining as u64).read_to_end(&mut drain);
uucore::io::RawWriter(stdout)
.write_all(&drain)
.inspect_err(handle_broken_pipe)
.map_err(write_err)?;
return Err(CatError::Io(unusable_errno()));
}
}
}

// splice is usable: from here on every error is a diagnosed error.
loop {
let mut remaining = match splice(reader, &pipe_wr, MAX_ROOTLESS_PIPE_SIZE) {
// End of input: splice handled the whole file.
Ok(0) => return Ok(()),
Ok(bytes_read) => bytes_read,
Err(errno) => return Err(CatError::Io(io::Error::from(errno))),
};
while remaining > 0 {
match splice(&pipe_rd, stdout, remaining) {
// No progress; stop splicing.
Ok(0) => return Ok(()),
Ok(bytes_written) => remaining -= bytes_written,
Err(errno) => return Err(write_err(io::Error::from(errno))),
}
}
}
}

/// Default splice failure: unusable, i.e. fold every pre-copy failure and
/// ENOSYS into EINVAL so callers test one errno (see
/// [`uucore::pipes::splice_unusable`]).
#[cfg(any(target_os = "linux", target_os = "android"))]
fn unusable_errno() -> io::Error {
io::Error::from_raw_os_error(rustix::io::Errno::INVAL.raw_os_error())
}

#[cfg_attr(any(target_os = "linux", target_os = "android"), inline(never))] // splice fast-path does not require this allocation
fn print_unbuffered<R: FdReadable>(
handle: &mut InputHandle<R>,
Expand All @@ -499,7 +599,8 @@ fn print_unbuffered<R: FdReadable>(
Ok(n) => {
stdout
.write_all(&buf[..n])
.inspect_err(handle_broken_pipe)?;
.inspect_err(handle_broken_pipe)
.map_err(write_err)?;
// cannot use rustix::io on Windows
// really bad workaround for unbuffered write <https://github.com/uutils/coreutils/issues/12188>
#[cfg(not(any(unix, target_os = "wasi")))]
Expand Down
10 changes: 8 additions & 2 deletions src/uu/tail/src/tail.rs
Original file line number Diff line number Diff line change
Expand Up @@ -586,8 +586,14 @@ fn print_target_section<
}
} else {
#[cfg(any(target_os = "linux", target_os = "android"))]
if uucore::pipes::splice_unbounded_auto(file, &mut stdout)?.is_err() {
io::copy(file, &mut stdout)?;
match uucore::pipes::splice_unbounded_auto(file, &mut stdout) {
// EINVAL means splice is unusable; copy with read/write.
Err(e) if uucore::pipes::splice_unusable(&e) => {
io::copy(file, &mut stdout)?;
}
// Real errors and success both need no further copying.
Ok(()) => {}
Err(e) => return Err(e.into()),
}
#[cfg(not(any(target_os = "linux", target_os = "android")))]
io::copy(file, &mut stdout)?;
Expand Down
16 changes: 12 additions & 4 deletions src/uu/tee/src/tee.rs
Original file line number Diff line number Diff line change
Expand Up @@ -155,10 +155,18 @@ impl MultiWriter {
macro_rules! splice_or_detach {
($pipe:expr, $writer:expr, $len:expr) => {
if let Err(e) = uucore::pipes::drain_pipe($pipe, $writer, $len) {
self.aborted |=
process_error(self.output_error_mode, e, $writer, &mut self.ignored_errors)
.is_err();
$writer.name.clear(); //mark as exited
// EINVAL means splice fell back to read/write: keep
// this writer, the data was still written.
if !uucore::pipes::splice_unusable(&e) {
self.aborted |= process_error(
self.output_error_mode,
e,
$writer,
&mut self.ignored_errors,
)
.is_err();
$writer.name.clear(); //mark as exited
}
}
};
}
Expand Down
16 changes: 11 additions & 5 deletions src/uucore/src/lib/features/buf_copy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,11 +15,17 @@ pub fn copy_fast(
src: &mut (impl std::io::Read + AsFd),
dest: &mut impl AsFd,
) -> std::io::Result<()> {
if crate::pipes::splice_unbounded_auto(src, dest)?.is_err() {
// fall back on writing "without buffering", or order of output would be wrong
// unrelated for cp /dev/stdin since cp does not have multiple input? <https://github.com/uutils/coreutils/issues/5186>
// RawWriter also removes io::copy's specialization e.g. copy_file_range which might use reflink
std::io::copy(src, &mut crate::io::RawWriter(dest))?;
match crate::pipes::splice_unbounded_auto(src, dest) {
// EINVAL means splice is unusable; fall back on read/write.
Err(e) if crate::pipes::splice_unusable(&e) => {
// fall back on writing "without buffering", or order of output would be wrong
// unrelated for cp /dev/stdin since cp does not have multiple input? <https://github.com/uutils/coreutils/issues/5186>
// RawWriter also removes io::copy's specialization e.g. copy_file_range which might use reflink
std::io::copy(src, &mut crate::io::RawWriter(dest))?;
}
// Real errors and success both need no further copying.
Ok(()) => {}
Err(e) => return Err(e),
}
Ok(())
}
Expand Down
78 changes: 52 additions & 26 deletions src/uucore/src/lib/features/pipes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,13 +17,24 @@ use std::{
pub const MAX_ROOTLESS_PIPE_SIZE: usize = 1024 * 1024;
const KERNEL_DEFAULT_PIPE_SIZE: usize = 64 * 1024;

/// A type allows to
/// - check that zero-copy succeed by ?.is_ok()
/// - check that zero-copy failed, but read/write fallback succeed by ?.is_err()
/// - catch the read/write fallback's error by ? or let Err(e)
/// Whether an error from the splice helpers means that splice is unusable
/// here, so the caller should fall back on read/write.
///
/// use rustix::io::Result for functions without read/write fallback
type PipeRes = std::io::Result<Result<(), ()>>;
/// The helpers in this module use `Err(EINVAL)` as that marker:
///
/// - `drain_pipe` fell back to read/write (the data was still written)
/// - `splice_unbounded_auto` could not splice anything
///
/// and they fold `ENOSYS` (kernel without the syscall) into it.
#[inline]
pub fn splice_unusable(err: &std::io::Error) -> bool {
err.raw_os_error() == Some(rustix::io::Errno::INVAL.raw_os_error())
}

#[inline]
fn splice_unusable_errno() -> std::io::Error {
std::io::Error::from_raw_os_error(rustix::io::Errno::INVAL.raw_os_error())
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1 fn is enough. not two.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think both of them is unnecessary since we can modify/check errno.


/// return pipe and try to extend its size
/// SIZE_REQUIRED should be true if you want to fail when changing pipe size failed
Expand Down Expand Up @@ -56,23 +67,34 @@ pub fn splice(source: &impl AsFd, target: &impl AsFd, len: usize) -> rustix::io:
}

/// splice `len` bytes from `pipe` into `dest`.
///
/// Returns `Err(EINVAL)` if splice turned out to be unusable: the data was
/// delivered by the read/write fallback instead, see [`splice_unusable`].
#[inline]
pub fn drain_pipe(pipe: &PipeReader, dest: &impl AsFd, len: usize) -> PipeRes {
pub fn drain_pipe(pipe: &PipeReader, dest: &impl AsFd, len: usize) -> std::io::Result<()> {
debug_assert!(len <= MAX_ROOTLESS_PIPE_SIZE, "unexpected RAM usage");
let mut remaining = len;
while remaining > 0 {
if let Ok(s) = splice(pipe, dest, remaining) {
remaining -= s;
} else {
// read/write fallback
// use read_to_end to make pipe empty for the case write failed
let mut drain = Vec::with_capacity(remaining);
pipe.take(remaining as u64).read_to_end(&mut drain)?;
RawWriter(&dest).write_all(&drain)?;
return Ok(Err(()));
match splice(pipe, dest, remaining) {
Ok(0) => {
// no progress; drain by hand

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I believe this match arm is misundertnding.

let mut drain = Vec::with_capacity(remaining);
pipe.take(remaining as u64).read_to_end(&mut drain)?;
RawWriter(&dest).write_all(&drain)?;
return Err(splice_unusable_errno());
}
Ok(s) => remaining -= s,
Err(_) => {
// read/write fallback
// use read_to_end to make pipe empty for the case write failed
let mut drain = Vec::with_capacity(remaining);
pipe.take(remaining as u64).read_to_end(&mut drain)?;
RawWriter(&dest).write_all(&drain)?;
return Err(splice_unusable_errno());
}
}
}
Ok(Ok(()))
Ok(())
}

/// check that source is FUSE
Expand All @@ -86,11 +108,14 @@ pub fn might_fuse(source: &impl AsFd) -> bool {
///
/// throughput is better than direct splice for the case one of in/output is pipe by unknown reason
/// This includes read ahead and optimization for stdout's pipe size
///
/// Returns `Err(EINVAL)` when nothing could be spliced, see [`splice_unusable`].
/// Errors while draining the intermediate pipe are real output errors.
#[inline]
pub fn splice_unbounded_auto(source: &impl AsFd, dest: &mut impl AsFd) -> PipeRes {
pub fn splice_unbounded_auto(source: &impl AsFd, dest: &mut impl AsFd) -> std::io::Result<()> {
static PIPE_CACHE: OnceLock<Option<(PipeReader, PipeWriter)>> = OnceLock::new();
let Some((pipe_rd, pipe_wr)) = PIPE_CACHE.get_or_init(|| pipe::<false>().ok()) else {
return Ok(Err(()));
return Err(splice_unusable_errno());
};

// fcntl for input would not improve throughput since
Expand All @@ -101,13 +126,11 @@ pub fn splice_unbounded_auto(source: &impl AsFd, dest: &mut impl AsFd) -> PipeRe
let _ = rustix::fs::fadvise(source, 0, None, rustix::fs::Advice::Sequential);
loop {
match splice(&source, &pipe_wr, MAX_ROOTLESS_PIPE_SIZE) {
Ok(0) => return Ok(Ok(())),
Ok(0) => return Ok(()),
Ok(n) => {
if drain_pipe(pipe_rd, dest, n)?.is_err() {
return Ok(Err(()));
}
drain_pipe(pipe_rd, dest, n)?;
}
Err(_) => return Ok(Err(())),
Err(_) => return Err(splice_unusable_errno()),
}
}
}
Expand Down Expand Up @@ -161,8 +184,11 @@ pub fn send_n_bytes(input: impl AsFd, target: impl AsFd, n: u64) -> std::io::Res
Ok(s) => {
n -= s as u64;
bytes_written += s as u64;
if drain_pipe(broker_r, &target, s)?.is_err() {
break false;
if let Err(e) = drain_pipe(broker_r, &target, s) {
if splice_unusable(&e) {
break false;
}
return Err(e);
}
}
_ => break false,
Expand Down
18 changes: 18 additions & 0 deletions tests/by-util/test_cat.rs
Original file line number Diff line number Diff line change
Expand Up @@ -856,6 +856,24 @@ fn test_write_error_handling() {
.stderr_contains("No space left on device");
}

/// Write errors must be diagnosed as "cat: write error: …" without naming
/// the input file.
#[test]
#[cfg(target_os = "linux")]
fn test_write_error_message() {
use std::fs::File;

let dev_full =
File::create("/dev/full").expect("Failed to open /dev/full - test must run on Linux");

new_ucmd!()
.pipe_in("test content that should cause write error to /dev/full")
.set_stdout(dev_full)
.fails()
.code_is(1)
.stderr_contains("cat: write error: No space left on device");
}

#[test]
#[cfg(target_os = "linux")]
fn test_version_help_dev_full() {
Expand Down
Loading