refactor: restructured streaming output with configurable buffering

- Introduced `StreamWriter` with configurable block and line buffering policies across builtins.
- Updated pipeline stages and compound commands to execute concurrently as tasks.
- Replaced Linux splice and local stream wrapping with generic `Write` streams and explicit flushing.
- Added regression and streaming smoke tests to verify concurrent output and prevent deadlocks.
This commit is contained in:
can1357
2026-08-20 03:14:09 +02:00
parent ee4a1d48b7
commit 12591dbde3
20 changed files with 820 additions and 335 deletions
+39 -148
View File
@@ -7,107 +7,19 @@ use std::os::unix::fs::FileTypeExt;
use std::{
ffi::OsString,
fs::{File, metadata},
io::{self, BufWriter, ErrorKind, Read, Write},
io::{self, ErrorKind, Read, Write},
path::Path,
};
use brush_core::{ShellExtensions, builtins::Registration, openfiles::OpenFile};
use clap::{Arg, ArgAction, ArgMatches, Command};
use memchr::memchr2;
use thiserror::Error;
use uucore::{display::Quotable, fast_inc::fast_inc_one};
use crate::host::{Host, Stdin, Utility, format_usage, matches_parser, util};
use brush_core::{ShellExtensions, builtins::Registration};
/// Linux splice support.
#[cfg(any(target_os = "linux", target_os = "android"))]
mod splice {
use std::io::{self, ErrorKind};
use std::os::fd::{AsFd, BorrowedFd};
use crate::host::{Host, Utility, format_usage, matches_parser, util};
use rustix::io::{read, write};
use uucore::pipes::{MAX_ROOTLESS_PIPE_SIZE, pipe, splice, splice_exact};
use super::{CatError, CatResult, FdReadable, InputHandle};
const BUF_SIZE: usize = 1024 * 16;
/// Moves input between real descriptors without copying through userspace.
///
/// `false` means the input reached EOF; `true` asks the caller to resume with
/// buffered copying because splice was unavailable for this descriptor pair.
#[inline]
pub(super) fn write_fast_using_splice<R: FdReadable>(
handle: &InputHandle<R>,
write_fd: BorrowedFd<'_>,
) -> CatResult<bool> {
let Some(read_fd) = handle.reader.try_borrow_as_fd() else {
return Ok(true);
};
if splice(&read_fd, &write_fd, MAX_ROOTLESS_PIPE_SIZE).is_ok() {
// fcntl improves throughput. It is harmless when stdout is not a pipe.
let _ = rustix::pipe::fcntl_setpipe_size(write_fd, MAX_ROOTLESS_PIPE_SIZE);
loop {
match splice(&read_fd, &write_fd, MAX_ROOTLESS_PIPE_SIZE) {
Ok(1..) => {},
Ok(0) => return Ok(false),
Err(error) if error.kind() == ErrorKind::BrokenPipe => {
return Err(CatError::BrokenPipe);
},
Err(_) => return Ok(true),
}
}
} else if let Ok((pipe_rd, pipe_wr)) = pipe() {
// Neither endpoint is a pipe, so broker through an intermediate pipe.
loop {
match splice(&read_fd, &pipe_wr, MAX_ROOTLESS_PIPE_SIZE) {
Ok(0) => return Ok(false),
Ok(n) => {
if let Err(error) = splice_exact(&pipe_rd, &write_fd, n) {
if error.kind() == ErrorKind::BrokenPipe {
return Err(CatError::BrokenPipe);
}
// Preserve bytes already moved into the broker pipe, then
// let the caller continue with buffered copying.
copy_exact(&pipe_rd, &write_fd, n)?;
return Ok(true);
}
},
Err(_) => return Ok(true),
}
}
} else {
Ok(true)
}
}
/// Moves exactly `num_bytes` bytes between the two descriptors.
fn copy_exact(
read_fd: &impl AsFd,
write_fd: &impl AsFd,
num_bytes: usize,
) -> io::Result<()> {
let mut left = num_bytes;
let mut buf = [0; BUF_SIZE];
while left > 0 {
let n = read(read_fd, &mut buf)?;
assert_ne!(n, 0, "unexpected end of pipe");
let mut written = 0;
while written < n {
match write(write_fd, &buf[written..n])? {
0 => unreachable!("fd should be writable"),
w => written += w,
}
}
left -= n;
}
Ok(())
}
}
// Allocate 32 digits for the line number. An estimate is that we can print
// about 1e8 lines/second, so 32 digits lasts for billions of universe lifetimes.
const LINE_NUMBER_BUF_SIZE: usize = 32;
struct LineNumber {
@@ -117,16 +29,23 @@ struct LineNumber {
num_end: usize,
}
// Manually incrementing the digits is significantly faster than formatting a
// `usize` each time. The buffer starts as " 1\t".
// Logic to store a string for the line number. Manually incrementing the value
// represented in a buffer like this is significantly faster than storing a
// `usize` and using the standard Rust formatting macros to format a `usize` to
// a string each time it's needed.
// Buffer is initialized to " 1\t" and incremented each time `increment` is
// called, using uucore's fast_inc function that operates on strings.
impl LineNumber {
fn new() -> Self {
let mut buf = [b'0'; LINE_NUMBER_BUF_SIZE];
let init_str = " 1\t";
let print_start = buf.len() - init_str.len();
let num_start = buf.len() - 2;
let num_end = buf.len() - 1;
buf[print_start..].copy_from_slice(init_str.as_bytes());
Self { buf, print_start, num_start, num_end }
}
@@ -232,25 +151,9 @@ struct OutputState {
one_blank_kept: bool,
}
trait FdReadable: Read {
#[cfg(any(target_os = "linux", target_os = "android"))]
fn try_borrow_as_fd(&self) -> Option<std::os::fd::BorrowedFd<'_>> {
None
}
}
impl FdReadable for File {
#[cfg(any(target_os = "linux", target_os = "android"))]
fn try_borrow_as_fd(&self) -> Option<std::os::fd::BorrowedFd<'_>> {
use std::os::fd::AsFd;
Some(self.as_fd())
}
}
impl FdReadable for &mut Stdin {}
/// An input stream and whether it is connected to an interactive terminal.
struct InputHandle<R: FdReadable> {
struct InputHandle<R: Read> {
reader: R,
is_interactive: bool,
}
@@ -327,7 +230,8 @@ impl Utility for Cat {
};
#[allow(clippy::unwrap_used, reason = "clap provides '-' by default")]
let files = self.matches.get_many::<OsString>(options::FILE).unwrap();
cat_files(files, &options, host)
let mut stdout = host.stdout_writer();
cat_files(files, &options, host, &mut stdout)
}
}
@@ -421,11 +325,11 @@ fn app() -> Command {
)
}
fn cat_handle<R: FdReadable>(
fn cat_handle<R: Read>(
handle: &mut InputHandle<R>,
options: &OutputOptions,
state: &mut OutputState,
stdout: &mut OpenFile,
stdout: &mut impl Write,
) -> CatResult<()> {
if options.can_write_fast() {
write_fast(handle, stdout)
@@ -439,13 +343,13 @@ fn cat_path(
options: &OutputOptions,
state: &mut OutputState,
host: &mut Host,
stdout: &mut impl Write,
) -> CatResult<()> {
// Resolve every operand at the boundary, but retain `path` for diagnostics.
let resolved = host.resolve(path);
match get_input_type(path, &resolved)? {
InputType::StdIn => {
let (stdin, stdout) = (&mut host.stdin, &mut host.stdout);
let mut handle = InputHandle { reader: stdin, is_interactive: false };
let mut handle = InputHandle { reader: &mut host.stdin, is_interactive: false };
cat_handle(&mut handle, options, state, stdout)
},
InputType::Directory => Err(CatError::IsDirectory),
@@ -454,12 +358,17 @@ fn cat_path(
_ => {
let file = File::open(resolved)?;
let mut handle = InputHandle { reader: file, is_interactive: false };
cat_handle(&mut handle, options, state, &mut host.stdout)
cat_handle(&mut handle, options, state, stdout)
},
}
}
fn cat_files<'a, I>(files: I, options: &OutputOptions, host: &mut Host) -> i32
fn cat_files<'a, I>(
files: I,
options: &OutputOptions,
host: &mut Host,
stdout: &mut impl Write,
) -> i32
where
I: IntoIterator<Item = &'a OsString>,
{
@@ -471,14 +380,15 @@ where
};
for path in files {
match cat_path(path, options, &mut state, host) {
match cat_path(path, options, &mut state, host, stdout) {
Ok(()) => {},
Err(CatError::BrokenPipe) => return host.exit_code(),
Err(error) => host.error(format!("{}: {error}", path.maybe_quote()), 1),
}
}
if state.skipped_carriage_return {
let _ = host.stdout.write_all(b"\r");
let _ = stdout.write_all(b"\r");
let _ = stdout.flush();
}
host.exit_code()
}
@@ -522,23 +432,10 @@ fn get_input_type(path: &OsString, resolved: &Path) -> CatResult<InputType> {
}
/// Writes a handle to stdout with no output transformation.
fn write_fast<R: FdReadable>(
fn write_fast<R: Read>(
handle: &mut InputHandle<R>,
stdout: &mut OpenFile,
stdout: &mut impl Write,
) -> CatResult<()> {
#[cfg(any(target_os = "linux", target_os = "android"))]
{
// Splice is safe only when the in-process output stream is backed by a
// real descriptor. Never substitute the host process's fd 1.
let fall_back = match stdout.try_borrow_as_fd() {
Ok(stdout_fd) => splice::write_fast_using_splice(handle, stdout_fd)?,
Err(_) => true,
};
if !fall_back {
return Ok(());
}
}
let mut buf = [0; 1024 * 64];
loop {
match handle.reader.read(&mut buf) {
@@ -553,15 +450,13 @@ fn write_fast<R: FdReadable>(
}
/// Outputs a handle line by line with the requested transformations.
fn write_lines<R: FdReadable>(
fn write_lines<R: Read>(
handle: &mut InputHandle<R>,
options: &OutputOptions,
state: &mut OutputState,
stdout: &mut OpenFile,
stdout: &mut impl Write,
) -> CatResult<()> {
let mut in_buf = [0; 1024 * 31];
// A 32K output buffer greatly improves performance.
let mut writer = BufWriter::with_capacity(32 * 1024, stdout);
loop {
let n = match handle.reader.read(&mut in_buf) {
@@ -574,23 +469,23 @@ fn write_lines<R: FdReadable>(
let mut pos = 0;
while pos < n {
if in_buf[pos] == b'\n' {
write_new_line(&mut writer, options, state, handle.is_interactive)?;
write_new_line(stdout, options, state, handle.is_interactive)?;
state.at_line_start = true;
pos += 1;
continue;
}
if state.skipped_carriage_return {
writer.write_all(b"\r")?;
stdout.write_all(b"\r")?;
state.skipped_carriage_return = false;
state.at_line_start = false;
}
state.one_blank_kept = false;
if state.at_line_start && options.number != NumberingMode::None {
state.line_number.write(&mut writer)?;
state.line_number.write(stdout)?;
state.line_number.increment();
}
let offset = write_end(&mut writer, &in_buf[pos..], options)?;
let offset = write_end(stdout, &in_buf[pos..], options)?;
if offset + pos == in_buf.len() {
state.at_line_start = false;
break;
@@ -599,17 +494,13 @@ fn write_lines<R: FdReadable>(
state.skipped_carriage_return = true;
} else {
assert_eq!(in_buf[pos + offset], b'\n');
write_end_of_line(
&mut writer,
options.end_of_line().as_bytes(),
handle.is_interactive,
)?;
write_end_of_line(stdout, options.end_of_line().as_bytes(), handle.is_interactive)?;
state.at_line_start = true;
}
pos += offset + 1;
}
// Flush before a pipe read can block so available output stays visible.
writer.flush()?;
stdout.flush()?;
}
Ok(())
}
+10 -8
View File
@@ -6,7 +6,7 @@ use std::{
cmp::Ordering,
ffi::{OsStr, OsString},
fs::{self, File},
io::{self, BufRead, BufReader, BufWriter, Read, Write},
io::{self, BufRead, BufReader, Read, Write},
path::Path,
};
@@ -148,7 +148,6 @@ fn compare(
usize::from(!opts.get_flag(options::COLUMN_1))
+ usize::from(!opts.get_flag(options::COLUMN_2)),
);
let mut writer = BufWriter::new(stdout);
let (mut ra, mut rb) = (Vec::new(), Vec::new());
let mut na = read_context(a, &mut ra, name1)?;
let mut nb = read_context(b, &mut rb, name2)?;
@@ -172,7 +171,7 @@ fn compare(
break;
}
if !opts.get_flag(options::COLUMN_1) {
writer.write_all(&ra).map_err(|e| format!("write error: {e}"))?;
stdout.write_all(&ra).map_err(|e| format!("write error: {e}"))?;
}
ra.clear();
na = read_context(a, &mut ra, name1)?;
@@ -183,7 +182,7 @@ fn compare(
break;
}
if !opts.get_flag(options::COLUMN_2) {
write_delimited(&mut writer, col2.as_bytes(), &rb)
write_delimited(&mut *stdout, col2.as_bytes(), &rb)
.map_err(|e| format!("write error: {e}"))?;
}
rb.clear();
@@ -197,7 +196,7 @@ fn compare(
break;
}
if !opts.get_flag(options::COLUMN_3) {
write_delimited(&mut writer, col3.as_bytes(), &ra)
write_delimited(&mut *stdout, col3.as_bytes(), &ra)
.map_err(|e| format!("write error: {e}"))?;
}
ra.clear();
@@ -213,10 +212,10 @@ fn compare(
}
if opts.get_flag(options::TOTAL) {
let ending = LineEnding::from_zero_flag(opts.get_flag(options::ZERO_TERMINATED));
write!(writer, "{n1}{delim}{n2}{delim}{n3}{delim}total{ending}")
write!(stdout, "{n1}{delim}{n2}{delim}{n3}{delim}total{ending}")
.map_err(|e| format!("write error: {e}"))?;
}
writer.flush().map_err(|e| format!("write error: {e}"))?;
stdout.flush().map_err(|e| format!("write error: {e}"))?;
if should_check && (c1.has_error || c2.has_error) {
if delayed_error {
let _ = writeln!(stderr, "comm: input is not in sorted order");
@@ -277,6 +276,9 @@ impl Utility for Comm {
files_identical(&path1, &path2).unwrap_or(false)
};
let ending = LineEnding::from_zero_flag(self.matches.get_flag(options::ZERO_TERMINATED));
// Taken before the `LineReader`s below hold `&mut host.stdin`; a
// method borrow of `host` would otherwise conflict with them.
let mut stdout = host.stdout_writer();
let opened: Result<_, (&OsStr, io::Error)> = if name1 == "-" {
open_file(name2, &path2, None, ending)
.map_err(|e| (name2.as_os_str(), e))
@@ -317,7 +319,7 @@ impl Utility for Comm {
delim,
&self.matches,
identical,
&mut host.stdout,
&mut stdout,
&mut host.stderr,
) {
Ok(true) => 0,
+2 -2
View File
@@ -5,7 +5,7 @@
use std::{
ffi::OsString,
fs::File,
io::{self, BufRead, BufReader, BufWriter, Read, Write},
io::{self, BufRead, BufReader, Read, Write},
};
use bstr::io::BufReadExt;
@@ -737,7 +737,7 @@ where
.collect::<Vec<_>>();
let mut stdin_read = false;
let mut failed = false;
let mut out = BufWriter::new(&mut host.stdout);
let mut out = host.stdout_writer();
for (filename, path) in inputs {
let result = if let Some(path) = path {
+2 -2
View File
@@ -763,7 +763,7 @@ use std::{
collections::HashMap,
ffi::OsString,
fs::File,
io::{BufRead, BufReader, BufWriter, Read, Write},
io::{BufRead, BufReader, Read, Write},
path::{Path, PathBuf},
sync::LazyLock,
};
@@ -1541,6 +1541,7 @@ fn date_main(host: &mut Host, matches: &ArgMatches) -> Result<(), DateError> {
let cancel = host.cancel_flag();
let mut had_error = false;
let mut debug_stderr = host.stderr_clone();
let mut stdout = host.stdout_writer();
let reader_stderr = host.stderr_clone();
let dates: Box<dyn Iterator<Item = _>> = match &settings.date_source {
DateSource::Human(input) => {
@@ -1732,7 +1733,6 @@ fn date_main(host: &mut Host, matches: &ArgMatches) -> Result<(), DateError> {
};
let format_string = make_format_string(&settings);
let mut stdout = BufWriter::new(&mut host.stdout);
// Format all the dates
let config = Config::new().custom(PosixCustom::new()).lenient(true);
+4 -4
View File
@@ -15,7 +15,7 @@ use std::{
time::{Duration, SystemTime, UNIX_EPOCH},
};
use brush_core::{ShellExtensions, builtins::Registration, openfiles::OpenFile};
use brush_core::{ShellExtensions, builtins::Registration};
use clap::{ArgAction, Parser, ValueEnum};
use globset::{GlobBuilder, GlobMatcher};
use pi_walker::CollectedEntry;
@@ -678,7 +678,7 @@ fn search(
}
let use_gitignore = !(no_ignore(&cli) || no_ignore_vcs(&cli));
let mut out = BufWriter::new(&mut host.stdout);
let mut out = host.stdout_writer();
let mut state = SearchState { matches: 0, had_error: false };
for search_path in &search_paths {
if cancelled.load(Ordering::Relaxed) || max_results.is_some_and(|max| state.matches >= max) {
@@ -743,8 +743,8 @@ fn try_search_fast(
search_paths: &[SearchPath],
config: &SearchConfig,
max_results: Option<usize>,
stdout: &mut OpenFile,
stderr: &mut OpenFile,
stdout: &mut impl Write,
stderr: &mut impl Write,
cancelled: &AtomicBool,
) -> io::Result<Option<SearchState>> {
if !can_use_fast_search(cli, config) {
+90 -3
View File
@@ -7,7 +7,7 @@ use std::{
borrow::Cow,
ffi::{OsStr, OsString},
fs::File,
io::{self, BufWriter, Read, Write},
io::{self, Read, Write},
path::{Path, PathBuf},
};
@@ -1274,7 +1274,7 @@ fn execute_search<M: Matcher>(
max_count: Option<u64>,
) -> i32 {
let mut searcher = build_searcher(cli, opts, max_count);
let mut out = BufWriter::new(host.stdout_clone());
let mut out = host.stdout_writer();
let mut any_match = false;
let mut had_error = false;
let mut processed_operand = false;
@@ -1523,14 +1523,101 @@ pub(crate) fn grep_builtin<SE: ShellExtensions>() -> Registration<SE> {
#[cfg(test)]
mod tests {
use std::{
io::{self, Read, Write},
sync::Arc,
};
use parking_lot::Mutex;
use super::*;
use crate::host::{Host, run_util};
use brush_core::openfiles;
use crate::host::{Host, run_caught, run_util};
struct SnapshottingStdin {
pos: usize,
snapped: bool,
stdout: Arc<Mutex<Option<Arc<Mutex<Vec<u8>>>>>>,
snapshot: Arc<Mutex<Vec<u8>>>,
}
const SNAPSHOT_INPUT: &[u8] = b"hit\nmiss\n";
impl Read for SnapshottingStdin {
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
if self.pos < SNAPSHOT_INPUT.len() {
let n = buf.len().min(SNAPSHOT_INPUT.len() - self.pos);
buf[..n].copy_from_slice(&SNAPSHOT_INPUT[self.pos..self.pos + n]);
self.pos += n;
return Ok(n);
}
// Input exhausted: grep is back asking for more. Whatever it has
// already flushed to stdout is what a live consumer would see now.
if !self.snapped {
let stdout = self.stdout.lock().clone().expect("stdout buffer is initialized");
*self.snapshot.lock() = stdout.lock().clone();
self.snapped = true;
}
Ok(0)
}
}
impl Write for SnapshottingStdin {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
Ok(buf.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
impl openfiles::Stream for SnapshottingStdin {
fn clone_box(&self) -> Box<dyn openfiles::Stream> {
Box::new(Self {
pos: self.pos,
snapped: self.snapped,
stdout: Arc::clone(&self.stdout),
snapshot: Arc::clone(&self.snapshot),
})
}
#[cfg(unix)]
fn try_clone_to_owned(&self) -> Result<std::os::fd::OwnedFd, brush_core::Error> {
Err(brush_core::error::ErrorKind::CannotConvertToNativeFd.into())
}
#[cfg(unix)]
fn try_borrow_as_fd(&self) -> Result<std::os::fd::BorrowedFd<'_>, brush_core::Error> {
Err(brush_core::error::ErrorKind::CannotConvertToNativeFd.into())
}
}
fn run(args: &[&str], stdin: &str) -> (i32, String, String) {
let (code, capture) = run_util::<Grep>(args, stdin, "/");
(code, capture.out(), capture.err())
}
#[test]
fn stdin_matches_are_visible_before_eof() {
let stdout = Arc::new(Mutex::new(None));
let snapshot = Arc::new(Mutex::new(Vec::new()));
let stdin = Box::new(SnapshottingStdin {
pos: 0,
snapped: false,
stdout: Arc::clone(&stdout),
snapshot: Arc::clone(&snapshot),
});
let (mut host, capture) = Host::for_test_with_stdin("grep", stdin, "/");
*stdout.lock() = Some(capture.stdout_buffer());
let parsed = Grep::try_parse_from(["grep", "hit", "-"]).unwrap();
assert_eq!(run_caught(parsed, &mut host), 0, "{}", capture.err());
// A regression re-buffering grep's output makes matches invisible until EOF.
assert_eq!(snapshot.lock().as_slice(), b"hit\n");
}
#[test]
fn max_count_and_no_match_statuses_are_gnu_compatible() {
let (code, out, err) = run(&["-m1", "hit"], "hit\nmiss\nhit\n");
+15 -12
View File
@@ -5,7 +5,7 @@
use std::{
ffi::OsString,
fs::File,
io::{self, BufWriter, Read, Seek, SeekFrom, Write},
io::{self, Read, Seek, SeekFrom, Write},
num::TryFromIntError,
path::PathBuf,
};
@@ -1166,9 +1166,8 @@ fn read_n_lines(
separator: u8,
) -> io::Result<u64> {
let mut reader = take_lines(input, n, separator);
let mut writer = BufWriter::with_capacity(BUF_SIZE, output);
let bytes_written = io::copy(&mut reader, &mut writer).map_err(wrap_in_stdout_error)?;
writer.flush().map_err(wrap_in_stdout_error)?;
let bytes_written = io::copy(&mut reader, output).map_err(wrap_in_stdout_error)?;
output.flush().map_err(wrap_in_stdout_error)?;
Ok(bytes_written)
}
@@ -1394,23 +1393,24 @@ impl Utility for Head {
let _ = out.write_all(b" <==\n");
*first = false;
}
let mut out = host.stdout_writer();
for file in &options.files {
let result = if file == "-" {
if print_headers {
print_header(&mut host.stdout, b"standard input", &mut first);
print_header(&mut out, b"standard input", &mut first);
}
let mut input = io::BufReader::with_capacity(BUF_SIZE, &mut host.stdin);
match options.mode {
Mode::FirstBytes(n) => read_n_bytes(&mut input, &mut host.stdout, n),
Mode::FirstBytes(n) => read_n_bytes(&mut input, &mut out, n),
Mode::AllButLastBytes(n) => {
read_but_last_n_bytes(&mut input, &mut host.stdout, n)
read_but_last_n_bytes(&mut input, &mut out, n)
},
Mode::FirstLines(n) => {
read_n_lines(&mut input, &mut host.stdout, n, options.line_ending.into())
read_n_lines(&mut input, &mut out, n, options.line_ending.into())
},
Mode::AllButLastLines(n) => read_but_last_n_lines(
&mut input,
&mut host.stdout,
&mut out,
n,
options.line_ending.into(),
),
@@ -1421,7 +1421,7 @@ impl Utility for Head {
// GNU prints the header before reporting the read error,
// and that header counts as produced output.
if print_headers {
print_header(&mut host.stdout, file.as_encoded_bytes(), &mut first);
print_header(&mut out, file.as_encoded_bytes(), &mut first);
}
host.error(format!("error reading {}: Is a directory", file.quote()), 1);
continue;
@@ -1434,9 +1434,9 @@ impl Utility for Head {
},
};
if print_headers {
print_header(&mut host.stdout, file.as_encoded_bytes(), &mut first);
print_header(&mut out, file.as_encoded_bytes(), &mut first);
}
head_file(&mut input, &mut host.stdout, &options)
head_file(&mut input, &mut out, &options)
};
if let Err(err) = result {
let name = if file == "-" {
@@ -1453,6 +1453,9 @@ impl Utility for Head {
}
}
}
if let Err(err) = out.flush() {
host.error(wrap_in_stdout_error(err), 1);
}
host.exit_code()
}
}
+317 -16
View File
@@ -30,7 +30,7 @@ use std::{
cell::Cell,
collections::HashMap,
ffi::OsString,
io::{self, Read, Write},
io::{self, BufWriter, LineWriter, Read, Write},
marker::PhantomData,
panic::{AssertUnwindSafe, catch_unwind},
path::{Path, PathBuf},
@@ -41,6 +41,8 @@ use std::{
},
};
use parking_lot::Mutex;
use brush_core::{
Error, ExecutionContext, ExecutionResult, ShellExtensions,
builtins::{self, Registration},
@@ -88,10 +90,15 @@ pub(crate) struct Host {
/// Standard input. Reads observe cancellation, so a blocked pipe read
/// returns EOF on abort instead of hanging the shell.
pub stdin: Stdin,
/// Standard output; the null device when fd 1 is closed.
/// Standard output; the null device when fd 1 is closed. Raw: utilities
/// with bulk output buffer it themselves via [`Host::stdout_writer`].
pub stdout: OpenFile,
/// Standard error; the null device when fd 2 is closed.
pub stderr: OpenFile,
/// Standard error, buffered with the destination-aware policy of
/// [`StreamWriter`]; the null device when fd 2 is closed. When fd 2 shares
/// fd 1's destination (`2>&1`, or the default capture pipe), this is the
/// same serialized writer [`Host::stdout_writer`] returns, so interleaving
/// follows write order exactly.
pub stderr: StreamWriter,
name: String,
cwd: PathBuf,
@@ -99,6 +106,9 @@ pub(crate) struct Host {
cancel: Arc<AtomicBool>,
exit_code: i32,
stdin_is_search_input: bool,
/// The shared stdout/stderr writer when both fds point at one
/// destination; `None` when they diverge.
merged_out: Option<Arc<Mutex<StreamWriter>>>,
}
struct CancelOnDrop(Arc<AtomicBool>);
@@ -195,9 +205,26 @@ impl Host {
self.stdout.clone()
}
/// Duplicates stderr, for utilities that hand a writer to a helper thread.
/// Duplicates stderr as a raw [`OpenFile`], for utilities that hand a
/// writer to a helper thread. Data pending in the buffered stderr (at most
/// one partial line) is not carried over.
pub fn stderr_clone(&self) -> OpenFile {
self.stderr.clone()
self.stderr.dup_file()
}
/// A buffered stdout with a flush policy chosen by the destination of
/// fd 1; see [`StdoutWriter`].
///
/// Utilities that emit output progressively — stream filters (`grep`,
/// `sed`, `cut`) and directory walkers (`ls`, `fd`) — must write through
/// this rather than a raw `BufWriter`, so their output is visible as it is
/// produced. Batch emitters whose output only exists once all input is
/// consumed (`sort`, `tac`, `seq`) may keep plain block buffering.
pub fn stdout_writer(&self) -> StreamWriter {
match &self.merged_out {
Some(shared) => StreamWriter::Shared(Arc::clone(shared)),
None => StreamWriter::new(self.stdout.clone()),
}
}
/// A launcher for child processes started by this utility.
@@ -215,7 +242,7 @@ impl Host {
.map(|(k, v)| (k.clone(), v.clone()))
.collect(),
),
stderr: self.stderr.clone(),
stderr: self.stderr.dup_file(),
}
}
@@ -260,6 +287,148 @@ impl Host {
}
}
/// Buffered writer for a utility's output streams, with a flush policy
/// matching the destination.
///
/// When the destination is a regular file (or the null device), nothing
/// observes the output until the utility exits, so writes are block-buffered
/// for throughput. Everywhere else — a pipe to the next pipeline stage, the
/// harness capture pipe behind the TUI's live tool output (a pipe fd wrapped
/// in `OpenFile::File`, hence the `fstat` in [`is_regular_file`] rather than
/// a variant match), or an in-memory stream — writes are line-buffered so
/// each completed line is visible as soon as it is produced rather than when
/// the utility exits.
///
/// Construct via [`Host::stdout_writer`]; [`StreamWriter::line`] and
/// [`StreamWriter::block`] force a policy for utilities with explicit
/// buffering flags (`rg --line-buffered`).
pub(crate) enum StreamWriter {
/// Block-buffered: flushed when full, on drop, and on explicit `flush`.
Block(BufWriter<OpenFile>),
/// Line-buffered: additionally flushed through the last newline of every
/// write.
Line(LineWriter<OpenFile>),
/// A serialized handle onto a writer shared by stdout and stderr, used
/// when fd 1 and fd 2 have the same destination (`2>&1`, or the default
/// capture pipe): one buffer means diagnostics and output interleave in
/// exactly the order they were written.
Shared(Arc<Mutex<StreamWriter>>),
}
impl StreamWriter {
const BLOCK_CAPACITY: usize = 64 * 1024;
const LINE_CAPACITY: usize = 16 * 1024;
/// Picks the policy for `file`: block for regular files, line otherwise.
pub fn new(file: OpenFile) -> Self {
if is_regular_file(&file) { Self::block(file) } else { Self::line(file) }
}
/// Forces line buffering regardless of destination.
pub fn line(file: OpenFile) -> Self {
Self::Line(LineWriter::with_capacity(Self::LINE_CAPACITY, file))
}
/// Forces block buffering regardless of destination.
pub fn block(file: OpenFile) -> Self {
Self::Block(BufWriter::with_capacity(Self::BLOCK_CAPACITY, file))
}
/// Duplicates the underlying descriptor as a raw [`OpenFile`], for
/// utilities that hand a writer to helper threads. Buffered data pending
/// in this writer (at most one partial line under the line policy) is not
/// carried over.
pub fn dup_file(&self) -> OpenFile {
match self {
Self::Block(w) => w.get_ref().clone(),
Self::Line(w) => w.get_ref().clone(),
Self::Shared(shared) => shared.lock().dup_file(),
}
}
/// Whether the destination is a terminal, mirroring
/// [`OpenFile::is_terminal`].
pub fn is_terminal(&self) -> bool {
match self {
Self::Block(w) => w.get_ref().is_terminal(),
Self::Line(w) => w.get_ref().is_terminal(),
Self::Shared(shared) => shared.lock().is_terminal(),
}
}
}
impl Write for StreamWriter {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
match self {
Self::Block(w) => w.write(buf),
Self::Line(w) => w.write(buf),
Self::Shared(shared) => shared.lock().write(buf),
}
}
fn flush(&mut self) -> io::Result<()> {
match self {
Self::Block(w) => w.flush(),
Self::Line(w) => w.flush(),
Self::Shared(shared) => shared.lock().flush(),
}
}
fn write_vectored(&mut self, bufs: &[io::IoSlice<'_>]) -> io::Result<usize> {
match self {
Self::Block(w) => w.write_vectored(bufs),
Self::Line(w) => w.write_vectored(bufs),
Self::Shared(shared) => shared.lock().write_vectored(bufs),
}
}
}
/// Whether writes to `file` land in a regular file, where output is only ever
/// observed after the utility exits.
///
/// A pipe wrapped in `std::fs::File` (how the shell hands the capture pipe to
/// a command) reports a fifo file type, and `metadata` on exotic handles can
/// fail outright; both classify as "not a regular file" and get line
/// buffering, the visibility-safe default.
pub(crate) fn is_regular_file(file: &OpenFile) -> bool {
match file {
OpenFile::File(f) => f.metadata().is_ok_and(|m| m.is_file()),
_ => false,
}
}
/// Whether two open files refer to the same non-seekable destination — the
/// `2>&1` case (and the harness default, where one capture pipe backs both
/// fds).
///
/// Matching is by `fstat` device+inode and deliberately excludes regular
/// files: `cmd >f 2>f` opens two descriptions with independent offsets, and
/// funneling them through one writer would change where the bytes land.
/// Pipes, fifos, terminals, and sockets have no offset, so a device+inode
/// match identifies the same object.
#[cfg(unix)]
fn same_destination(a: &OpenFile, b: &OpenFile) -> bool {
use std::os::unix::fs::MetadataExt;
fn id(file: &OpenFile) -> Option<(u64, u64)> {
let fd = file.try_borrow_as_fd().ok()?;
let dup = fd.try_clone_to_owned().ok()?;
let meta = std::fs::File::from(dup).metadata().ok()?;
if meta.file_type().is_file() {
return None;
}
Some((meta.dev(), meta.ino()))
}
match (id(a), id(b)) {
(Some(a), Some(b)) => a == b,
_ => false,
}
}
#[cfg(not(unix))]
fn same_destination(_a: &OpenFile, _b: &OpenFile) -> bool {
false
}
/// A shell-faithful launcher for child processes started by a utility builtin.
///
/// Carries the three things a child must inherit from the *shell* rather than
@@ -639,7 +808,7 @@ async fn run_utility<U: Utility, SE: ShellExtensions>(
/// the long-lived host process. With `panic = "unwind"` the panic unwinds to
/// here, where it becomes a non-zero exit plus a concise note on the command's
/// own stderr.
fn run_caught<U: Utility>(parsed: U, host: &mut Host) -> i32 {
pub(crate) fn run_caught<U: Utility>(parsed: U, host: &mut Host) -> i32 {
struct Guard;
impl Drop for Guard {
fn drop(&mut self) {
@@ -698,20 +867,32 @@ fn build_host<SE: ShellExtensions>(
// `Stdin::read` must observe the very same flag or it never wakes.
let cancel = Arc::new(AtomicBool::new(false));
let stdout = or_null(context.try_fd(OpenFiles::STDOUT_FD))?;
let stderr_file = or_null(context.try_fd(OpenFiles::STDERR_FD))?;
// `2>&1` (and the default capture pipe): one shared writer keeps
// diagnostics and output in exact write order.
let (merged_out, stderr) = if same_destination(&stdout, &stderr_file) {
let shared = Arc::new(Mutex::new(StreamWriter::new(stderr_file)));
(Some(Arc::clone(&shared)), StreamWriter::Shared(shared))
} else {
(None, StreamWriter::new(stderr_file))
};
Ok(Host {
stdin: Stdin {
file: or_null(stdin)?,
fd: stdin_fd,
cancel: Arc::clone(&cancel),
},
stdout: or_null(context.try_fd(OpenFiles::STDOUT_FD))?,
stderr: or_null(context.try_fd(OpenFiles::STDERR_FD))?,
stdout,
stderr,
name: invoked,
cwd: context.shell.working_dir().to_path_buf(),
env,
cancel,
exit_code: 0,
stdin_is_search_input,
merged_out,
})
}
@@ -808,8 +989,8 @@ mod testing {
use parking_lot::Mutex;
use super::{
Arc, AtomicBool, HashMap, Host, OpenFile, OsString, PathBuf, Read, Stdin, Utility, Write, io,
openfiles, run_caught,
Arc, AtomicBool, HashMap, Host, OpenFile, OsString, PathBuf, Read, Stdin, StreamWriter,
Utility, Write, io, openfiles, run_caught,
};
/// Captured in-memory output from [`Host::for_test`].
@@ -824,6 +1005,11 @@ mod testing {
self.stdout.lock().clone()
}
/// Shared stdout buffer for tests that must observe output mid-run.
pub(crate) fn stdout_buffer(&self) -> Arc<Mutex<Vec<u8>>> {
Arc::clone(&self.stdout)
}
/// Raw bytes the utility wrote to stderr.
pub fn stderr(&self) -> Vec<u8> {
self.stderr.lock().clone()
@@ -849,6 +1035,19 @@ mod testing {
name: &str,
stdin: impl Into<Vec<u8>>,
cwd: impl Into<PathBuf>,
) -> (Self, Capture) {
Self::for_test_with_stdin(
name,
Box::new(MemStream::reader(stdin.into())),
cwd,
)
}
/// Builds a host backed by an arbitrary in-memory stdin stream.
pub(crate) fn for_test_with_stdin(
name: &str,
stdin: Box<dyn openfiles::Stream>,
cwd: impl Into<PathBuf>,
) -> (Self, Capture) {
let capture = Capture {
stdout: Arc::new(Mutex::new(Vec::new())),
@@ -857,22 +1056,23 @@ mod testing {
let cancel = Arc::new(AtomicBool::new(false));
let host = Self {
stdin: Stdin {
file: OpenFile::Stream(Box::new(MemStream::reader(stdin.into()))),
file: OpenFile::Stream(stdin),
fd: None,
cancel: Arc::clone(&cancel),
},
stdout: OpenFile::Stream(Box::new(MemStream::writer(Arc::clone(
&capture.stdout,
)))),
stderr: OpenFile::Stream(Box::new(MemStream::writer(Arc::clone(
&capture.stderr,
)))),
stderr: StreamWriter::new(OpenFile::Stream(Box::new(
MemStream::writer(Arc::clone(&capture.stderr)),
))),
name: name.to_string(),
cwd: cwd.into(),
env: HashMap::new(),
cancel,
exit_code: 0,
stdin_is_search_input: false,
merged_out: None,
};
(host, capture)
}
@@ -984,6 +1184,107 @@ mod testing {
Err(brush_core::error::ErrorKind::CannotConvertToNativeFd.into())
}
}
mod stdout_policy {
use parking_lot::Mutex;
use super::MemStream;
use crate::host::{Arc, OpenFile, StreamWriter, Write};
/// Contract: on a non-file destination, a completed line is visible to
/// the consumer before any explicit flush; a partial line is held back.
#[test]
fn line_policy_flushes_completed_lines_immediately() {
let buf = Arc::new(Mutex::new(Vec::new()));
let stream = OpenFile::Stream(Box::new(MemStream::writer(Arc::clone(&buf))));
let mut out = StreamWriter::new(stream);
assert!(matches!(out, StreamWriter::Line(_)));
out.write_all(b"hit\n").unwrap();
assert_eq!(buf.lock().as_slice(), b"hit\n");
out.write_all(b"partial").unwrap();
assert_eq!(buf.lock().as_slice(), b"hit\n");
out.flush().unwrap();
assert_eq!(buf.lock().as_slice(), b"hit\npartial");
}
/// Contract: a regular-file destination stays block-buffered — bytes
/// reach the file only on flush, not per line.
#[test]
fn regular_file_gets_block_buffering() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("out.txt");
let file = std::fs::File::create(&path).unwrap();
let mut out = StreamWriter::new(OpenFile::File(file));
assert!(matches!(out, StreamWriter::Block(_)));
out.write_all(b"hit\n").unwrap();
assert_eq!(std::fs::read(&path).unwrap(), b"");
out.flush().unwrap();
assert_eq!(std::fs::read(&path).unwrap(), b"hit\n");
}
/// Contract: the shell hands commands their stdout as a pipe fd wrapped
/// in `std::fs::File`; that must classify as line-buffered, or live tool
/// output stalls until the utility exits.
#[cfg(unix)]
#[test]
fn pipe_wrapped_as_file_gets_line_buffering() {
let (reader, writer) = std::io::pipe().unwrap();
let file = std::fs::File::from(std::os::fd::OwnedFd::from(writer));
assert!(matches!(StreamWriter::new(OpenFile::File(file)), StreamWriter::Line(_)));
drop(reader);
}
/// Contract for `2>&1`: two handles onto one shared writer interleave
/// in exact write order — diagnostics land where they were emitted
/// relative to output, not where a second buffer happened to flush.
#[test]
fn shared_handles_preserve_write_order() {
let buf = Arc::new(Mutex::new(Vec::new()));
let inner =
StreamWriter::line(OpenFile::Stream(Box::new(MemStream::writer(Arc::clone(&buf)))));
let shared = Arc::new(Mutex::new(inner));
let mut out = StreamWriter::Shared(Arc::clone(&shared));
let mut err = StreamWriter::Shared(shared);
writeln!(out, "out 1").unwrap();
writeln!(err, "err 1").unwrap();
writeln!(out, "out 2").unwrap();
assert_eq!(buf.lock().as_slice(), b"out 1\nerr 1\nout 2\n");
}
/// Contract: `2>&1` over a pipe is detected (same object, no offset),
/// while distinct pipes and regular files — which have independent
/// offsets under `>f 2>f` — are not merged.
#[cfg(unix)]
#[test]
fn same_destination_detects_dup_pipes_only() {
use crate::host::same_destination;
let (reader, writer) = std::io::pipe().unwrap();
let dup = writer.try_clone().unwrap();
let a = OpenFile::File(std::fs::File::from(std::os::fd::OwnedFd::from(writer)));
let b = OpenFile::File(std::fs::File::from(std::os::fd::OwnedFd::from(dup)));
assert!(same_destination(&a, &b));
let (reader2, writer2) = std::io::pipe().unwrap();
let c = OpenFile::File(std::fs::File::from(std::os::fd::OwnedFd::from(writer2)));
assert!(!same_destination(&a, &c));
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("out.txt");
let f1 = OpenFile::File(std::fs::File::create(&path).unwrap());
let f2 = OpenFile::File(std::fs::File::create(&path).unwrap());
assert!(!same_destination(&f1, &f2));
drop((reader, reader2));
}
}
}
#[cfg(test)]
+8 -8
View File
@@ -839,11 +839,10 @@ mod output {
}
}
/// Runs `f` with buffered standard output.
/// Runs `f` with standard output.
pub fn with_stdout<T>(stdout: &mut dyn Write, f: impl FnOnce(&mut dyn Write) -> T) -> T {
let mut out = io::BufWriter::new(stdout);
let res = f(&mut out);
let _ = out.flush();
let res = f(stdout);
let _ = stdout.flush();
res
}
@@ -928,7 +927,8 @@ impl Utility for Jq {
let _runtime = RuntimeGuard::install(host);
color::set(!cli.in_place && cli.color_if(|| stdout_is_terminal));
let result = real_main(&cli, host);
let mut stdout = host.stdout_writer();
let result = real_main(&cli, host, &mut stdout);
if let Some(code) = filter::take_halt() {
return code;
}
@@ -1036,7 +1036,7 @@ mod color {
}
}
fn real_main(cli: &Cli, host: &mut Host) -> Result<i32, Error> {
fn real_main(cli: &Cli, host: &mut Host, stdout: &mut dyn Write) -> Result<i32, Error> {
if let Some(test_files) = &cli.run_tests {
return Ok(match test_files.last() {
Some(file) => {
@@ -1072,7 +1072,7 @@ fn real_main(cli: &Cli, host: &mut Host) -> Result<i32, Error> {
let last = if cli.files.is_empty() {
let inputs = read::buffered(cli, io::BufReader::new(&mut host.stdin));
output::with_stdout(&mut host.stdout, |out| {
output::with_stdout(stdout, |out| {
filter::run(cli, &filter, ctx, inputs, |v| output::print(out, cli, &v))
})?
} else {
@@ -1104,7 +1104,7 @@ fn real_main(cli: &Cli, host: &mut Host) -> Result<i32, Error> {
tmp.persist(path).map_err(Error::Persist)?;
std::fs::set_permissions(path, perms)?;
} else {
last = output::with_stdout(&mut host.stdout, |out| {
last = output::with_stdout(stdout, |out| {
filter::run(cli, &filter, ctx.clone(), inputs, |v| output::print(out, cli, &v))
})?;
}
+18 -19
View File
@@ -10,7 +10,7 @@ use std::{
cmp::Reverse,
ffi::{OsStr, OsString},
fs::{self, DirEntry, FileType, Metadata, ReadDir},
io::{BufWriter, ErrorKind, Write},
io::{ErrorKind, Write},
ops::RangeInclusive,
path::{Path, PathBuf},
rc::Rc,
@@ -37,7 +37,7 @@ use uucore::{
version_cmp::version_cmp,
};
use crate::host::{Host, Utility, format_usage, matches_parser, os_bytes_lossy, util};
use crate::host::{Host, StreamWriter, Utility, format_usage, matches_parser, os_bytes_lossy, util};
mod colors {
//! Color handling for the `ls` builtin.
@@ -1996,7 +1996,7 @@ mod dired {
use std::{
fmt,
io::{self, BufWriter, Write},
io::{self, Write},
};
/// `dired` Module Documentation
@@ -2068,7 +2068,7 @@ pub fn calculate_dired(
(start, end)
}
pub fn indent<W: Write>(out: &mut BufWriter<W>) -> io::Result<()> {
pub fn indent<W: Write>(out: &mut W) -> io::Result<()> {
write!(out, " ")?;
Ok(())
}
@@ -2085,7 +2085,7 @@ pub fn calculate_subdired(dired: &mut DiredOutput, path_len: usize) {
pub fn print_dired_output<W: Write>(
config: &Config,
dired: &DiredOutput,
out: &mut BufWriter<W>,
out: &mut W,
) -> io::Result<()> {
out.flush()?;
if !dired.dired_positions.is_empty() {
@@ -2102,7 +2102,7 @@ pub fn print_dired_output<W: Write>(
/// Helper function to print positions with a given prefix.
fn print_positions<W: Write>(
out: &mut BufWriter<W>,
out: &mut W,
prefix: &str,
positions: &[BytePosition],
) -> io::Result<()> {
@@ -2318,17 +2318,17 @@ use std::os::unix::fs::{FileTypeExt, MetadataExt};
#[cfg(windows)]
use std::os::windows::fs::MetadataExt;
/// Show the directory name in the case where several arguments are given to ls
use std::{borrow::Cow, iter};
use std::{
borrow::Cow,
cell::LazyCell,
ffi::{OsStr, OsString},
fmt::Write as FmtWrite,
fs::{self, DirEntry, FileType, Metadata},
io::{BufWriter, Write},
io::Write,
iter,
sync::LazyLock,
time::SystemTime,
};
use brush_core::openfiles::OpenFile;
use ansi_width::ansi_width;
use glob::MatchOptions;
@@ -2442,9 +2442,9 @@ enum SizeOrDeviceId {
/// dir1: <- This as well
/// file11
/// ```
pub fn show_dir_name(
pub fn show_dir_name<W: Write>(
path_data: &PathData,
out: &mut BufWriter<OpenFile>,
out: &mut W,
config: &Config,
) -> std::io::Result<()> {
let escaped_name = escape_dir_name_with_locale(path_data.path().as_os_str(), config);
@@ -2761,11 +2761,11 @@ pub fn display_items(
Ok(())
}
fn display_grid(
fn display_grid<W: Write>(
names: impl Iterator<Item = OsString>,
width: u16,
direction: Direction,
out: &mut BufWriter<OpenFile>,
out: &mut W,
quoted: bool,
tab_size: usize,
) -> std::io::Result<()> {
@@ -3155,8 +3155,7 @@ fn display_item_name(
DisplayItemName { displayed: name, dired_name_len }
}
/// This writes to the [`BufWriter`] `state.out` a single string of the output
/// of `ls -l`.
/// This writes to `state.out` a single string of the output of `ls -l`.
///
/// It writes the following keys, in order:
/// * `inode` ([`display_inode`], config-optional)
@@ -4738,7 +4737,7 @@ type DirData = (PathBuf, bool);
// A struct to encapsulate state that is passed around from `list` functions.
#[cfg_attr(not(unix), allow(dead_code))]
struct ListState<'a> {
out: BufWriter<OpenFile>,
out: StreamWriter,
style_manager: Option<StyleManager<'a>>,
// TODO: More benchmarking with different use cases is required here.
// From experiments, BTreeMap may be faster than HashMap, especially as the
@@ -4769,7 +4768,7 @@ pub fn list(locs: Vec<&Path>, config: &Config, stdout: OpenFile) -> std::io::Res
let now = SystemTime::now();
let mut state = ListState {
out: BufWriter::new(stdout),
out: StreamWriter::new(stdout),
style_manager: config
.color
.as_ref()
@@ -5124,10 +5123,10 @@ fn get_metadata_with_deref_opt(path: &Path, dereference: bool) -> std::io::Resul
}
}
fn write_total(
fn write_total<W: Write>(
items: &[PathData],
config: &Config,
out: &mut BufWriter<OpenFile>,
out: &mut W,
) -> std::io::Result<usize> {
let mut total_size = 0;
for item in items {
+6 -26
View File
@@ -12,11 +12,10 @@
use std::{
ffi::{OsStr, OsString},
fs::File,
io::{self, BufWriter, Read, Write},
io::{self, Read, Write},
path::{Path, PathBuf},
};
use brush_core::openfiles::OpenFile;
use clap::{ArgAction, Parser, ValueEnum};
use grep_cli::DecompressionReaderBuilder;
use grep_matcher::{Captures, LineTerminator, Matcher};
@@ -26,7 +25,7 @@ use grep_regex::{RegexMatcher, RegexMatcherBuilder};
use grep_searcher::{
BinaryDetection, Encoding, Searcher, SearcherBuilder, Sink, SinkContext, SinkFinish, SinkMatch,
};
use crate::host::{Host, Utility};
use crate::host::{Host, StreamWriter, Utility};
use ignore::{
Match,
@@ -498,27 +497,6 @@ enum CompiledMatcher {
Pcre(PcreMatcher),
}
enum RgOutput {
Buffered(BufWriter<OpenFile>),
Direct(OpenFile),
}
impl Write for RgOutput {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
match self {
Self::Buffered(output) => output.write(bytes),
Self::Direct(output) => output.write(bytes),
}
}
fn flush(&mut self) -> io::Result<()> {
match self {
Self::Buffered(output) => output.flush(),
Self::Direct(output) => output.flush(),
}
}
}
struct SearchOptions {
line_number: bool,
column: bool,
@@ -1840,9 +1818,11 @@ impl Utility for Rg {
return 2;
}
let mut out = if cli.line_buffered && !cli.no_line_buffered {
RgOutput::Direct(host.stdout_clone())
StreamWriter::line(host.stdout_clone())
} else if cli.no_line_buffered {
StreamWriter::block(host.stdout_clone())
} else {
RgOutput::Buffered(BufWriter::new(host.stdout_clone()))
host.stdout_writer()
};
let (patterns, mut paths) = match resolve_patterns(host, &cli) {
Ok(resolved) => resolved,
+75 -20
View File
@@ -5628,6 +5628,7 @@ struct MmapOutput {
/// All other output is buffered and writen via BufWriter.
pub struct OutputBuffer {
out: BufWriter<Box<dyn OutputWrite + 'static>>, // Where to write
line_buffered: bool, // Flush completed lines for non-file stdout
#[cfg(unix)]
max_pending_write: usize, /* Max bytes to keep before
* flushing */
@@ -5654,9 +5655,10 @@ const MAX_PENDING_WRITE_NON_FILE: usize = 64 * 1024;
impl OutputBuffer {
#[cfg(not(unix))]
pub fn new(w: Box<dyn OutputWrite + 'static>) -> Self {
pub fn new(w: Box<dyn OutputWrite + 'static>, line_buffered: bool) -> Self {
Self {
out: BufWriter::new(w),
line_buffered,
pending_newline: false,
#[cfg(test)]
low_level_flushes: 0,
@@ -5664,12 +5666,15 @@ impl OutputBuffer {
}
#[cfg(unix)]
pub fn new(w: Box<dyn OutputWrite + 'static>) -> Self {
// The writer is not fd-backed, so regular-file output detection is gone;
// always bound pending data by the pipe-sized limit.
pub fn new(w: Box<dyn OutputWrite + 'static>, line_buffered: bool) -> Self {
Self {
out: BufWriter::new(w),
max_pending_write: MAX_PENDING_WRITE_NON_FILE,
line_buffered,
max_pending_write: if line_buffered {
MAX_PENDING_WRITE_NON_FILE
} else {
usize::MAX
},
mmap_chunk: None,
pending_newline: false,
#[cfg(test)]
@@ -5701,7 +5706,33 @@ impl OutputBuffer {
};
let mut reader = BufReader::new(file);
io::copy(&mut reader, &mut self.out)?;
if self.line_buffered {
let mut buf = [0; 8 * 1024];
loop {
let len = reader.read(&mut buf)?;
if len == 0 {
break;
}
self.out.write_all(&buf[..len])?;
if buf[..len].contains(&b'\n') {
self.flush_completed_line()?;
}
}
} else {
io::copy(&mut reader, &mut self.out)?;
}
Ok(())
}
/// Flush output through a completed line when writing to a non-file stdout.
fn flush_completed_line(&mut self) -> io::Result<()> {
if self.line_buffered {
#[cfg(test)]
{
self.low_level_flushes += 1;
}
self.out.flush()?;
}
Ok(())
}
}
@@ -5740,6 +5771,7 @@ impl OutputBuffer {
self.flush_mmap(WriteRange::Complete)?;
self.out.write_all(b"\n")?;
self.pending_newline = false;
self.flush_completed_line()?;
}
match &new_chunk.content {
@@ -5789,6 +5821,13 @@ impl OutputBuffer {
self.pending_newline = !has_newline;
},
}
if self.line_buffered && new_chunk.is_newline_terminated() {
// Mmap output reaches the BufWriter only here; file-backed mmap
// output is block-buffered, so this cannot affect its fast path.
self.flush_mmap(WriteRange::Complete)?;
self.flush_completed_line()?;
}
Ok(())
}
@@ -5819,6 +5858,7 @@ impl OutputBuffer {
self.flush_mmap(WriteRange::Complete)?;
self.out.write_all(b"\n")?;
self.pending_newline = false;
self.flush_completed_line()?;
}
Ok(())
}
@@ -5841,6 +5881,7 @@ impl OutputBuffer {
if self.pending_newline {
self.out.write_all(b"\n")?;
self.pending_newline = false;
self.flush_completed_line()?;
}
match &chunk.content {
@@ -5850,6 +5891,9 @@ impl OutputBuffer {
self.out.write_all(b"\n")?;
}
self.pending_newline = !has_newline;
if *has_newline {
self.flush_completed_line()?;
}
Ok(())
},
}
@@ -5860,6 +5904,7 @@ impl OutputBuffer {
if self.pending_newline {
self.out.write_all(b"\n")?;
self.pending_newline = false;
self.flush_completed_line()?;
}
Ok(())
}
@@ -5904,7 +5949,7 @@ mod tests {
let tmp = NamedTempFile::new()?;
{
let file = tmp.reopen()?;
let mut out = OutputBuffer::new(Box::new(file));
let mut out = OutputBuffer::new(Box::new(file), false);
out.write_str("foo\n")?;
out.write_str("bar\n")?;
out.flush()?;
@@ -5939,7 +5984,7 @@ mod tests {
let output = NamedTempFile::new()?;
let output_path = output.path().to_path_buf();
let out_file = std::fs::File::create(&output_path)?;
let mut out = OutputBuffer::new(Box::new(Box::new(out_file)));
let mut out = OutputBuffer::new(Box::new(Box::new(out_file)), false);
// Drain reader → writer
while let Some(chunk) = reader.get_line()? {
@@ -5974,7 +6019,7 @@ mod tests {
let output = NamedTempFile::new()?;
let output_path = output.path().to_path_buf();
let out_file = File::create(&output_path)?;
let mut out = OutputBuffer::new(Box::new(out_file));
let mut out = OutputBuffer::new(Box::new(out_file), false);
// Read the first mmap line ("zero\n") and write it
if let Some(chunk) = reader.get_line()? {
@@ -6028,7 +6073,7 @@ mod tests {
let out_file = File::create(&output_path)?;
// Wrap it in your OutputBuffer and run the loop:
let mut out = OutputBuffer::new(Box::new(out_file));
let mut out = OutputBuffer::new(Box::new(out_file), false);
let mut nline = 0;
while let Some(chunk) = reader.get_line()? {
out.write_chunk(&chunk)?;
@@ -6067,7 +6112,7 @@ mod tests {
let out_file = File::create(&output_path)?;
// Wrap it in your OutputBuffer and run the loop:
let mut out = OutputBuffer::new(Box::new(out_file));
let mut out = OutputBuffer::new(Box::new(out_file), false);
let mut nline = 0;
while let Some(chunk) = reader.get_line()? {
out.write_chunk(&chunk)?;
@@ -6102,7 +6147,7 @@ mod tests {
let out_file = File::create(&output_path)?;
// Wrap it in your OutputBuffer and run the loop:
let mut out = OutputBuffer::new(Box::new(out_file));
let mut out = OutputBuffer::new(Box::new(out_file), false);
let mut nline = 0;
while let Some(chunk) = reader.get_line()? {
out.write_chunk(&chunk)?;
@@ -6137,7 +6182,7 @@ mod tests {
let out_file = File::create(&output_path)?;
// Wrap it in your OutputBuffer and run the loop:
let mut out = OutputBuffer::new(Box::new(out_file));
let mut out = OutputBuffer::new(Box::new(out_file), false);
let mut nline = 0;
while let Some(chunk) = reader.get_line()? {
out.write_chunk(&chunk)?;
@@ -6441,6 +6486,7 @@ mod tests {
let file = tempfile().unwrap();
let buf = OutputBuffer {
out: BufWriter::new(Box::new(file.try_clone().unwrap())),
line_buffered: false,
#[cfg(unix)]
max_pending_write: 8,
#[cfg(unix)]
@@ -7333,9 +7379,14 @@ use tempfile::NamedTempFile;
use uucore::display::Quotable;
use brush_core::openfiles::OpenFile;
use crate::sed::error_handling::{IoContext, SedError, SedResult};
use crate::sed::{command::ProcessingContext, fast_io::OutputBuffer};
use crate::{
host::is_regular_file,
sed::{
command::ProcessingContext,
error_handling::{IoContext, SedError, SedResult},
fast_io::OutputBuffer,
},
};
/// Context for in-place editing
pub struct InPlace {
@@ -7353,9 +7404,10 @@ impl InPlace {
/// Depending on its settings it may or may not perform in-place
/// editing, backup the original file, or follow symlinks.
pub fn new_with_stdout(context: ProcessingContext, stdout: OpenFile) -> Self {
let line_buffered = !is_regular_file(&stdout);
Self {
stdout: stdout.clone(),
output: OutputBuffer::new(Box::new(stdout.clone())),
output: OutputBuffer::new(Box::new(stdout.clone()), line_buffered),
in_place: context.in_place,
in_place_suffix: context.in_place_suffix,
follow_symlinks: context.follow_symlinks,
@@ -7389,7 +7441,8 @@ impl InPlace {
/// to the context settings.
fn begin_resolved(&mut self, file_name: &Path) -> SedResult<&mut OutputBuffer> {
if !self.in_place {
self.output = OutputBuffer::new(Box::new(self.stdout.clone()));
self.output =
OutputBuffer::new(Box::new(self.stdout.clone()), !is_regular_file(&self.stdout));
return Ok(&mut self.output);
}
@@ -7418,8 +7471,10 @@ impl InPlace {
fs::set_permissions(temp_file.path(), perms)?;
}
let output =
OutputBuffer::new(Box::new(temp_file.reopen().expect("reopening NamedTempFile")));
let output = OutputBuffer::new(
Box::new(temp_file.reopen().expect("reopening NamedTempFile")),
false,
);
self.output = output;
self.temp_file = Some(temp_file);
self.original_path = Some(file_name.to_path_buf());
+23 -23
View File
@@ -1330,7 +1330,7 @@ mod follow {
use std::{
collections::{HashMap, hash_map::Keys},
fs::{File, Metadata},
io::{BufRead, BufReader, BufWriter, Write},
io::{BufRead, BufReader, Write},
path::{Path, PathBuf},
};
@@ -1480,8 +1480,7 @@ mod follow {
self.header_printer.print(display_name.as_str(), writer);
}
let mut writer = BufWriter::new(writer);
chunks.print(&mut writer).map_err(crate::tail::map_output_error)?;
chunks.print(writer).map_err(crate::tail::map_output_error)?;
writer.flush().map_err(crate::tail::map_output_error)?;
self.last.replace(path.to_owned());
@@ -1565,7 +1564,7 @@ mod follow {
use brush_core::openfiles::OpenFile;
use crate::{
host::Host,
host::{Host, StreamWriter},
tail::{
TailError,
TailResult,
@@ -1656,7 +1655,7 @@ mod follow {
pub files: FileHandling,
pub pid: platform::Pid,
pub stdout: OpenFile,
pub stdout: StreamWriter,
pub stderr: OpenFile,
pub cancel: Arc<AtomicBool>,
}
@@ -1668,7 +1667,7 @@ mod follow {
use_polling: bool,
files: FileHandling,
pid: platform::Pid,
stdout: OpenFile,
stdout: StreamWriter,
stderr: OpenFile,
cancel: Arc<AtomicBool>,
) -> Self {
@@ -1694,7 +1693,7 @@ mod follow {
pub fn from(
settings: &Settings,
stdout: OpenFile,
stdout: StreamWriter,
stderr: OpenFile,
cancel: Arc<AtomicBool>,
) -> Self {
@@ -2807,7 +2806,7 @@ use std::{
cmp::Ordering,
ffi::OsString,
fs::File,
io::{self, BufReader, BufWriter, ErrorKind, Read, Seek, SeekFrom, Write},
io::{self, BufReader, ErrorKind, Read, Seek, SeekFrom, Write},
path::{Path, PathBuf},
};
@@ -3075,6 +3074,7 @@ fn reverse_main(settings: &Settings, all_lines: bool, host: &mut Host) -> TailRe
unreachable!("-r with -c is rejected before dispatch");
};
let (signum, sep) = (*signum, *sep);
let mut stdout = host.stdout_writer();
let mut printer = HeaderPrinter::new(settings.verbose, true);
for input in &settings.inputs {
let path = match input.kind() {
@@ -3087,7 +3087,7 @@ fn reverse_main(settings: &Settings, all_lines: bool, host: &mut Host) -> TailRe
if let Some(path) = path {
if path.is_dir() {
host.fail(1);
printer.print_input(input, &mut host.stdout);
printer.print_input(input, &mut stdout);
let _ = writeln!(
host.stderr,
"tail: error reading '{}': Is a directory",
@@ -3097,7 +3097,7 @@ fn reverse_main(settings: &Settings, all_lines: bool, host: &mut Host) -> TailRe
}
match File::open(path) {
Ok(mut file) => {
printer.print_input(input, &mut host.stdout);
printer.print_input(input, &mut stdout);
file.read_to_end(&mut data)?;
},
Err(error) if error.kind() == ErrorKind::NotFound => {
@@ -3120,11 +3120,12 @@ fn reverse_main(settings: &Settings, all_lines: bool, host: &mut Host) -> TailRe
},
}
} else {
printer.print_input(input, &mut host.stdout);
printer.print_input(input, &mut stdout);
host.stdin.read_to_end(&mut data)?;
}
write_reversed_lines(&data, signum, sep, all_lines, &mut host.stdout)?;
write_reversed_lines(&data, signum, sep, all_lines, &mut stdout)?;
}
stdout.flush()?;
Ok(())
}
@@ -3164,7 +3165,6 @@ fn write_reversed_lines(
},
}
};
let mut writer = BufWriter::new(writer);
for segment in keep.iter().rev() {
writer.write_all(segment)?;
}
@@ -3192,7 +3192,7 @@ fn uu_tail(settings: &Settings, host: &mut Host) -> TailResult<()> {
let mut printer = HeaderPrinter::new(settings.verbose, true);
let mut observer = Observer::from(
settings,
host.stdout_clone(),
host.stdout_writer(),
host.stderr_clone(),
host.cancel_flag(),
);
@@ -3219,6 +3219,7 @@ fn uu_tail(settings: &Settings, host: &mut Host) -> TailResult<()> {
},
}
}
observer.stdout.flush()?;
if settings.follow.is_some() {
/*
@@ -3574,16 +3575,15 @@ fn unbounded_tail<T: Read>(
settings: &Settings,
writer: &mut impl Write,
) -> io::Result<()> {
let mut writer = BufWriter::new(writer);
match &settings.mode {
FilterMode::Lines(Signum::Negative(count), sep) => {
let mut chunks = chunks::LinesChunkBuffer::new(*sep, *count);
chunks.fill(reader)?;
chunks.write(&mut writer)?;
chunks.write(&mut *writer)?;
},
FilterMode::Lines(Signum::PlusZero | Signum::Positive(1), _) => {
io::copy(reader, &mut writer)?;
io::copy(reader, &mut *writer)?;
},
FilterMode::Lines(Signum::Positive(count), sep) => {
let mut num_skip = *count - 1;
@@ -3597,22 +3597,22 @@ fn unbounded_tail<T: Read>(
}
}
if chunk.has_data() {
chunk.write_lines(&mut writer, num_skip as usize)?;
io::copy(reader, &mut writer)?;
chunk.write_lines(&mut *writer, num_skip as usize)?;
io::copy(reader, &mut *writer)?;
}
},
FilterMode::Bytes(Signum::Negative(count)) => {
let mut chunks = chunks::BytesChunkBuffer::new(*count);
chunks.fill(reader)?;
chunks.print(&mut writer)?;
chunks.print(&mut *writer)?;
},
FilterMode::Lines(Signum::MinusZero, sep) => {
let mut chunks = chunks::LinesChunkBuffer::new(*sep, 0);
chunks.fill(reader)?;
chunks.write(&mut writer)?;
chunks.write(&mut *writer)?;
},
FilterMode::Bytes(Signum::PlusZero | Signum::Positive(1)) => {
io::copy(reader, &mut writer)?;
io::copy(reader, &mut *writer)?;
},
FilterMode::Bytes(Signum::Positive(count)) => {
let mut num_skip = *count - 1;
@@ -3635,7 +3635,7 @@ fn unbounded_tail<T: Read>(
}
}
io::copy(reader, &mut writer)?;
io::copy(reader, &mut *writer)?;
},
_ => {},
}
+6 -7
View File
@@ -847,17 +847,16 @@ fn run_uniq(matches: &ArgMatches, host: &mut Host) -> PortResult<()> {
})
.transpose()?;
// Writer first: `stdout_writer` method-borrows `host`, which must not
// overlap the `&mut host.stdin` held by the reader.
let writer: Box<dyn Write + '_> = match output_file {
Some(file) => Box::new(BufWriter::with_capacity(OUTPUT_BUFFER_CAPACITY, file)),
None => Box::new(host.stdout_writer()),
};
let reader: Box<dyn BufRead + '_> = match input_file {
Some(file) => Box::new(BufReader::new(file)),
None => Box::new(BufReader::new(&mut host.stdin)),
};
let writer: Box<dyn Write + '_> = match output_file {
Some(file) => Box::new(BufWriter::with_capacity(OUTPUT_BUFFER_CAPACITY, file)),
None => Box::new(BufWriter::with_capacity(
OUTPUT_BUFFER_CAPACITY,
&mut host.stdout,
)),
};
uniq.write_uniq(reader, writer)
}
+35
View File
@@ -46,6 +46,41 @@ pub fn block_range_at(options: BlockRangeOptions) -> Result<Option<BlockRange>>
.map_err(|error| Error::from_reason(error.to_string()))
}
#[napi(object)]
pub struct NodeSpan {
/// 1-indexed inclusive first line of the node.
pub start_line: u32,
/// 1-indexed inclusive last content line of the node.
pub end_line: u32,
/// Tree-sitter grammar node kind (e.g. `attribute_item`, `function_item`).
pub kind: String,
}
impl From<pi_ast::block::NodeSpan> for NodeSpan {
fn from(value: pi_ast::block::NodeSpan) -> Self {
Self { start_line: value.start_line, end_line: value.end_line, kind: value.kind }
}
}
/// Named-node chain containing `options.line`, innermost-first, excluding the
/// whole-file root.
///
/// Single-line nodes beginning on the line (attributes, decorators) come
/// first, followed by every enclosing construct. ERROR/MISSING recovery nodes
/// are skipped. Returns `null` when the language is unrecognized, the line is
/// out of range / blank, or the source fails to parse entirely.
#[napi]
pub fn node_chain_at(options: BlockRangeOptions) -> Result<Option<Vec<NodeSpan>>> {
pi_ast::block::node_chain_at(pi_ast::block::BlockRangeOptions {
code: options.code,
lang: options.lang,
path: options.path,
line: options.line,
})
.map(|chain| chain.map(|spans| spans.into_iter().map(Into::into).collect()))
.map_err(|error| Error::from_reason(error.to_string()))
}
#[napi(object)]
pub struct LineRange {
/// 1-indexed inclusive first visible line.
+4 -9
View File
@@ -47,7 +47,8 @@ struct Arena {
/// Stored as `u16` units purely for the 2-alignment UTF-16 fills need;
/// UTF-8 fills reinterpret the same bytes at alignment 1.
buf: UnsafeCell<[u16; SCRATCH_LEN / 2]>,
/// Bytes handed out. Fills bump it; drops roll it back (see [`Self::release`]).
/// Bytes handed out. Fills bump it; drops roll it back (see
/// [`Self::release`]).
offset: Cell<usize>,
/// Live scratch-backed guards. Hitting zero resets `offset`, so a non-LIFO
/// drop order leaks at most until the last guard goes away.
@@ -179,10 +180,7 @@ pub fn utf16(value: JsString<'_>) -> Result<Utf16> {
napi::check_status!(status, "Failed to read JavaScript string")?;
if written < avail - 1 {
arena.commit(start, written * 2);
return Ok(Utf16(TextRepr::Scratch {
ptr: NonNull::new(ptr).unwrap(),
len: written,
}));
return Ok(Utf16(TextRepr::Scratch { ptr: NonNull::new(ptr).unwrap(), len: written }));
}
}
@@ -236,10 +234,7 @@ pub fn utf8(value: JsString<'_>) -> Result<Utf8> {
return Err(Error::new(Status::InvalidArg, error.to_string()));
}
arena.commit(start, written);
return Ok(Utf8(TextRepr::Scratch {
ptr: NonNull::new(ptr).unwrap(),
len: written,
}));
return Ok(Utf8(TextRepr::Scratch { ptr: NonNull::new(ptr).unwrap(), len: written }));
}
}
+81
View File
@@ -0,0 +1,81 @@
//! Pipeline streaming contracts: builtin/compound/function pipeline stages
//! run concurrently and deliver output chunks while the command is still
//! running, instead of buffering until exit.
#![cfg(unix)]
use std::time::{Duration, Instant};
use pi_shell::{
cancel::CancelToken,
shell::{ShellExecuteOptions, execute_shell},
};
async fn run_collecting(command: &str) -> (Option<Duration>, Duration, String) {
let (tx, rx) = flume::unbounded::<String>();
let start = Instant::now();
let collector = tokio::spawn(async move {
let mut first_hit: Option<Duration> = None;
let mut output = String::new();
while let Ok(chunk) = rx.recv_async().await {
if first_hit.is_none() && chunk.contains("hit") {
first_hit = Some(start.elapsed());
}
output.push_str(&chunk);
}
(first_hit, output)
});
let result = execute_shell(
ShellExecuteOptions {
command: command.to_string(),
timeout_ms: Some(30_000),
..Default::default()
},
Some(tx),
CancelToken::new(None),
)
.await
.expect("shell execution");
let total = start.elapsed();
assert_eq!(result.exit_code, Some(0), "command failed: {command}");
let (first_hit, output) = collector.await.expect("collector task");
(first_hit, total, output)
}
/// A pipeline stage's output must reach the chunk callback while the pipeline
/// is still running. Regression: compound and function stages used to execute
/// inline during pipeline spawn, so downstream stages (and the consumer) saw
/// nothing until the stage exited; utility builtins used to hold stdout in an
/// exit-flushed `BufWriter`.
#[tokio::test]
async fn pipeline_stages_stream_before_exit() {
for command in [
"echo hit; sleep 1",
"{ echo hit; sleep 1; } | cat",
"{ echo hit; sleep 1; } | grep .",
"f() { echo hit; sleep 1; }; f | cat",
] {
let (first_hit, total, _) = run_collecting(command).await;
let first_hit = first_hit.expect("saw the first output chunk");
// The producer sleeps 1s after emitting the first line; seeing that
// line only in the back half of the run means it was buffered.
assert!(
first_hit < total / 2,
"`{command}`: first chunk at {first_hit:?}, finished at {total:?}: buffered until exit"
);
}
}
/// Regression: a compound stage that outfills the connecting pipe's buffer
/// used to deadlock the whole shell — the stage ran inline during pipeline
/// spawn, blocking on a pipe whose reader had not been spawned yet.
#[tokio::test]
async fn compound_stage_larger_than_pipe_buffer_does_not_deadlock() {
let (_, _, output) = run_collecting("{ seq 1 200000; echo hit; } | head -n 1").await;
// The merged channel may also carry seq's EPIPE diagnostic once head
// closes the pipe; the contract is that head emitted its line and the
// pipeline terminated.
assert_eq!(output.lines().next(), Some("1"));
}
+46 -18
View File
@@ -20,7 +20,7 @@ use crate::{
interp::{self, Execute, ExternalCommandInfo, ExternalCommandOutputMarkers, ProcessGroupPolicy},
openfiles::{self, OpenFiles},
pathsearch, processes,
results::ExecutionSpawnResult,
results::{ExecutionSpawnResult, ExecutionWaitResult},
sys, trace_categories, traps, variables,
};
@@ -512,32 +512,60 @@ impl<'a, SE: extensions::ShellExtensions> SimpleCommand<'a, SE> {
self,
func_registration: functions::Registration,
) -> Result<ExecutionSpawnResult, error::Error> {
let mut shell = self.shell;
let mut params = self.params;
params.disable_command_output_marking();
let last_arg = Self::take_last_arg(&self.args);
let cmd_context = ExecutionContext {
shell: &mut shell,
command_name: self.command_name,
params,
};
match self.shell {
// The function runs in an owned subshell (a pipeline stage or
// async job): execute it as a task, mirroring
// `execute_via_builtin_in_owned_shell`, so all pipeline stages run
// concurrently and its output streams to the next stage as it is
// produced.
ShellForCommand::OwnedShell { target, .. } => {
let mut shell = *target;
let command_name = self.command_name;
let args = self.args;
let join_handle = tokio::spawn(async move {
let cmd_context = ExecutionContext { shell: &mut shell, command_name, params };
let result =
invoke_shell_function(func_registration, cmd_context, &args[1..]).await;
// Strip the function name off args.
let result = invoke_shell_function(func_registration, cmd_context, &self.args[1..]).await;
// $_ is reset *after* the function body runs; see the
// parent-shell path below.
shell.update_last_arg_variable(last_arg);
// $_ is reset *after* the function body runs, to the last argument of
// the invocation (or the function name itself if zero args). Any
// mutations made inside the body are overwritten — this matches bash,
// where the caller observes only the invocation's last argument.
shell.update_last_arg_variable(last_arg);
match result?.wait().await? {
ExecutionWaitResult::Completed(result) => Ok(result),
ExecutionWaitResult::Stopped(_) => Ok(ExecutionResult::stopped()),
}
});
Ok(ExecutionSpawnResult::StartedTask(join_handle))
},
mut shell => {
let cmd_context = ExecutionContext {
shell: &mut shell,
command_name: self.command_name,
params,
};
if let Some(post_execute) = self.post_execute {
let _ = post_execute(&mut shell);
// Strip the function name off args.
let result = invoke_shell_function(func_registration, cmd_context, &self.args[1..]).await;
// $_ is reset *after* the function body runs, to the last argument of
// the invocation (or the function name itself if zero args). Any
// mutations made inside the body are overwritten — this matches bash,
// where the caller observes only the invocation's last argument.
shell.update_last_arg_variable(last_arg);
if let Some(post_execute) = self.post_execute {
let _ = post_execute(&mut shell);
}
result
},
}
result
}
fn execute_via_external(self, path: &Path) -> Result<ExecutionSpawnResult, error::Error> {
+34 -10
View File
@@ -928,17 +928,41 @@ impl<SE: extensions::ShellExtensions> ExecuteInPipeline<SE> for ast::Command {
Self::Simple(simple) => simple.execute_in_pipeline(pipeline_context, params).await,
Self::Compound(compound, redirects) => {
params.disable_command_output_marking();
// Set up any additional redirects.
if let Some(redirects) = redirects {
for redirect in &redirects.0 {
setup_redirect(&mut pipeline_context.shell, &mut params, redirect).await?;
}
}
Ok(compound
.execute(&mut pipeline_context.shell, &params)
.await?
.into())
// Each stage of a multi-command pipeline runs in its own
// subshell (`pipeline_context.shell` is already an owned
// clone). Execute compound stages as concurrent tasks rather
// than inline: inline execution serializes the pipeline —
// downstream stages are not even spawned until this stage
// completes — and deadlocks outright once a stage fills the
// connecting pipe's buffer with no reader running.
let in_pipeline = pipeline_context.in_pipeline;
match pipeline_context.shell {
commands::ShellForCommand::OwnedShell { target, .. } if in_pipeline => {
let mut shell = *target;
let compound = compound.clone();
let redirects = redirects.clone();
let mut params = params;
Ok(ExecutionSpawnResult::StartedTask(tokio::spawn(async move {
if let Some(redirects) = &redirects {
for redirect in &redirects.0 {
setup_redirect(&mut shell, &mut params, redirect).await?;
}
}
compound.execute(&mut shell, &params).await
})))
},
mut shell => {
// Set up any additional redirects.
if let Some(redirects) = redirects {
for redirect in &redirects.0 {
setup_redirect(&mut shell, &mut params, redirect).await?;
}
}
Ok(compound.execute(&mut shell, &params).await?.into())
},
}
},
Self::Function(func) => {
params.disable_command_output_marking();
+5
View File
@@ -7,6 +7,11 @@
- Added `ClaudeV3`/`ClaudeV47`/`ClaudeV5` encodings to `countTokens`: a Rust rewrite of [ctok](https://github.com/sanderland/ctok) by Sander Land (MIT), reconstructing Anthropic's `count_tokens` offline. Counts are exact on ctok's ~3.4M-response measurement corpora; the port is validated against 493 Python-ctok reference fixtures covering all three families. The pipeline is byte-level throughout — markers occupy one byte, normalization borrows text no rule touches, ASCII and ideographs skip the Unicode tables, and pieces are matched with one Aho-Corasick transition per byte instead of a per-position vocabulary descent — which counts English prose at 64 MiB/s, markdown at 73 MiB/s, source code at 35 MiB/s and CJK at 49 MiB/s per core: 1.5× (CJK, already cheap per byte) to 5.5× (prose, markdown, digits) a straightforward character-level implementation of the same model, which is held to byte-for-byte identical counts across 2.4M randomized differential comparisons.
- Added zstd-embedded exact content tokenizers for Qwen 3.5+/3.6+/3.8, DeepSeek V3/V4/R1, Kimi K2/K3, and GLM-5 alongside the rebuilt OpenAI o200k/cl100k and Claude reconstructions. `countTokens` now reads JavaScript strings through a reusable UTF-16 buffer, so native counting does not allocate a UTF-8 temporary.
### Changed
- Shell builtin utilities now stream their output. Utilities that emit progressively (`grep`, `rg`, `sed`, `cat`, `head`, `tail`, `cut`, `date`, `uniq`, `comm`, `jq`, `ls`, `fd`) write through a destination-aware buffer: line-buffered to pipes/terminals so each completed line reaches the TUI's live tool output (or the next pipeline stage) as it is produced, block-buffered to regular files for throughput. Previously they held everything in an exit-flushed 8–32 KiB `BufWriter`. `rg --line-buffered`/`--no-line-buffered` still force a policy. Builtin stderr goes through the same policy, and when fd 1 and fd 2 share a destination (`2>&1`, or the default merged capture pipe) both streams funnel through one serialized writer, so diagnostics interleave with output in exact write order.
- Compound (`{ …; }`, `(…)`) and shell-function pipeline stages now run concurrently with the rest of the pipeline, like builtin and external stages already did. They previously executed inline while the pipeline was being spawned, which delayed all downstream output until the stage exited and deadlocked the shell when a stage produced more than a pipe buffer with no reader running (e.g. `{ seq 1 200000; } | head -n 1`).
## [17.3.8] - 2026-08-19
### Changed