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
29 changes: 19 additions & 10 deletions src/uu/tee/src/tee.rs
Original file line number Diff line number Diff line change
Expand Up @@ -127,23 +127,32 @@ fn copy(mut input: impl Read, mut output: impl Write) -> Result<usize> {
};
let mut buffer = [0u8; FIRST_BUF_SIZE];
let mut len = 0;
match input.read(&mut buffer) {
Ok(0) => return Ok(0),
Ok(bytes_count) => {
output.write_all(&buffer[0..bytes_count])?;
len = bytes_count;
if bytes_count < FIRST_BUF_SIZE {
// Read with the stack buffer until we observe a full-sized read, which
// is a hint that the input is big enough to benefit from a larger
// buffer. A *short* read does not imply EOF: `read(2)` only returns 0
// at end-of-file, and a pipeline whose writer pauses between writes
// will commonly return less than the buffer size per call. We must
// therefore keep looping until we either see a full-sized read (and
// upgrade to the larger buffer) or `read` returns 0.
loop {
match input.read(&mut buffer) {
Ok(0) => return Ok(len), // end of file
Ok(bytes_count) => {
output.write_all(&buffer[..bytes_count])?;
// flush the buffer to comply with POSIX requirement that
// `tee` does not buffer the input.
output.flush()?;
return Ok(len);
len += bytes_count;
if bytes_count == FIRST_BUF_SIZE {

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 == is too strict and difficult to switch to code path for large file.

break;
}
}
Err(e) if e.kind() == ErrorKind::Interrupted => {}
Err(e) => return Err(e),
}
Err(e) if e.kind() == ErrorKind::Interrupted => (),
Err(e) => return Err(e),
}

// but optimize buffer size also for large file
// Optimize buffer size for large files.
let mut buffer = vec![0u8; 4 * FIRST_BUF_SIZE]; //stack array makes code path for smaller file slower
loop {
match input.read(&mut buffer) {
Expand Down
58 changes: 58 additions & 0 deletions tests/by-util/test_tee.rs
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,64 @@ fn test_tee_output_not_buffered() {
handle.join().unwrap();
}

#[test]
fn test_tee_continues_after_short_read() {
// Regression test: `tee` must keep reading until EOF even when the
// first `read(2)` returns fewer bytes than its internal buffer. This
// happens in any pipeline where the upstream writer pauses between
// writes (e.g. a slow producer, a `sleep` in a shell pipeline, or a
// service emitting log lines in bursts). Treating a short read as
// end-of-file caused tee to exit prematurely and downstream producers
// to die with SIGPIPE on their next write.
//
// Run in a separate thread so that the test fails via timeout rather
// than hanging if a regression reintroduces the bug in a form where
// tee blocks on read instead of exiting early.
let handle = std::thread::spawn(move || {
let (at, mut ucmd) = at_and_ucmd!();
let file_out = "tee_short_read_out";

let mut child = ucmd
.arg(file_out)
.set_stdin(Stdio::piped())
.set_stdout(Stdio::piped())
.run_no_wait();

// First chunk — deliberately much smaller than tee's internal
// buffer so that `read(2)` returns a short count.
child.write_in(b"first\n");
assert_eq!(&child.stdout_exact_bytes(6), b"first\n");

// Give a buggy implementation time to exit before we try to
// write again.
child.delay(50);

// Second chunk. A correctly-implemented tee is still reading
// from stdin; a buggy one has already exited and this write
// will either fail with EPIPE or never reach the output file.
child.write_in(b"second\n");
assert_eq!(&child.stdout_exact_bytes(7), b"second\n");

// `wait` closes stdin for us before waiting on the child.
child.wait().unwrap().success();

assert_eq!(at.read(file_out), "first\nsecond\n");
});

for _ in 0..500 {
std::thread::sleep(Duration::from_millis(10));
if handle.is_finished() {
break;
}
}

assert!(
handle.is_finished(),
"tee did not complete within the timeout"
);
handle.join().unwrap();
}

#[cfg(target_os = "linux")]
mod linux_only {
use uutests::util::{AtPath, CmdResult, UCommand};
Expand Down
Loading