From 97f9d493abe0ff03fa15ba71e01a8ead739708b7 Mon Sep 17 00:00:00 2001 From: Anthony DePasquale Date: Thu, 13 Aug 2026 10:42:15 +0200 Subject: [PATCH 1/2] sort: propagate output write errors --- src/uu/sort/src/ext_sort/threaded.rs | 53 ++++++++++++-- src/uu/sort/src/merge.rs | 101 ++++++++++++++++++++++++--- tests/by-util/test_sort.rs | 16 +++++ 3 files changed, 157 insertions(+), 13 deletions(-) diff --git a/src/uu/sort/src/ext_sort/threaded.rs b/src/uu/sort/src/ext_sort/threaded.rs index 7dd089d0fe8..c4ef133d1f0 100644 --- a/src/uu/sort/src/ext_sort/threaded.rs +++ b/src/uu/sort/src/ext_sort/threaded.rs @@ -293,13 +293,56 @@ fn write( separator: u8, ) -> UResult { let mut tmp_file = I::create(file, compress_prog)?; - write_lines(chunk.lines(), tmp_file.as_write(), separator); - tmp_file.finished_writing() + let write_result = write_lines(chunk.lines(), tmp_file.as_write(), separator); + let finish_result = tmp_file.finished_writing(); + write_result?; + finish_result } -fn write_lines(lines: &[Line], writer: &mut T, separator: u8) { +fn write_lines(lines: &[Line], writer: &mut T, separator: u8) -> std::io::Result<()> { for s in lines { - writer.write_all(s.line).unwrap(); - writer.write_all(&[separator]).unwrap(); + writer.write_all(s.line)?; + writer.write_all(&[separator])?; + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use std::io::{self, Write}; + + use super::write_lines; + use crate::Line; + + struct FailAfterFirstWrite { + writes: usize, + } + + impl Write for FailAfterFirstWrite { + fn write(&mut self, buf: &[u8]) -> io::Result { + if self.writes == 0 { + self.writes += 1; + Ok(buf.len()) + } else { + Err(io::Error::other("write failed")) + } + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + } + + #[test] + fn write_lines_propagates_write_errors() { + let lines = [Line { + line: b"line", + index: 0, + }]; + let mut writer = FailAfterFirstWrite { writes: 0 }; + + let error = write_lines(&lines, &mut writer, b'\n').unwrap_err(); + + assert_eq!(error.kind(), io::ErrorKind::Other); } } diff --git a/src/uu/sort/src/merge.rs b/src/uu/sort/src/merge.rs index 8f0b5bd54ae..bf85e49fc5f 100644 --- a/src/uu/sort/src/merge.rs +++ b/src/uu/sort/src/merge.rs @@ -25,7 +25,9 @@ use std::{ }; use compare::Compare; +use uucore::display::Quotable; use uucore::error::{FromIo, UResult}; +use uucore::translate; use crate::{ GlobalSettings, Output, SortError, @@ -140,8 +142,10 @@ pub fn merge_with_file_limit< let mut tmp_file = Tmp::create(tmp_dir.next_file()?, settings.compress_prog.as_deref())?; - merger.write_all_to(settings, tmp_file.as_write())?; - temporary_files.push(tmp_file.finished_writing()?); + let write_result = merger.write_all_to(settings, tmp_file.as_write()); + let finish_result = tmp_file.finished_writing(); + write_result?; + temporary_files.push(finish_result?); } } // Merge any remaining files that didn't get merged in a full batch above. @@ -151,8 +155,10 @@ pub fn merge_with_file_limit< let mut tmp_file = Tmp::create(tmp_dir.next_file()?, settings.compress_prog.as_deref())?; - merger.write_all_to(settings, tmp_file.as_write())?; - temporary_files.push(tmp_file.finished_writing()?); + let write_result = merger.write_all_to(settings, tmp_file.as_write()); + let finish_result = tmp_file.finished_writing(); + write_result?; + temporary_files.push(finish_result?); } merge_with_file_limit::<_, _, Tmp>( temporary_files @@ -304,8 +310,14 @@ struct FileMerger<'a> { impl FileMerger<'_> { /// Write the merged contents to the output file. fn write_all(self, settings: &GlobalSettings, output: Output) -> UResult<()> { + let output_name = output + .as_output_name() + .unwrap_or(OsStr::new("standard output")) + .to_owned(); + let ctx = || translate!("sort-error-write-failed", "output" => output_name.maybe_quote()); let mut out = output.into_write(); - self.write_all_to(settings, &mut out) + self.write_all_to(settings, &mut out)?; + flush_writer(&mut out).map_err_context(ctx) } fn write_all_to(mut self, settings: &GlobalSettings, out: &mut impl Write) -> UResult<()> { @@ -447,6 +459,10 @@ fn check_child_success(mut child: Child, program: &str) -> UResult<()> { } } +fn flush_writer(writer: &mut impl Write) -> std::io::Result<()> { + writer.flush() +} + /// A temporary file that can be written to. pub trait WriteableTmpFile: Sized { type Closed: ClosedTmpFile; @@ -493,7 +509,8 @@ impl WriteableTmpFile for WriteablePlainTmpFile { }) } - fn finished_writing(self) -> UResult { + fn finished_writing(mut self) -> UResult { + flush_writer(&mut self.file)?; Ok(ClosedPlainTmpFile { path: self.path }) } @@ -565,9 +582,12 @@ impl WriteableTmpFile for WriteableCompressedTmpFile { }) } - fn finished_writing(self) -> UResult { + fn finished_writing(mut self) -> UResult { + let flush_result = flush_writer(&mut self.child_stdin); drop(self.child_stdin); - check_child_success(self.child, &self.compress_prog)?; + let child_result = check_child_success(self.child, &self.compress_prog); + flush_result?; + child_result?; Ok(ClosedCompressedTmpFile { path: self.path, compress_prog: self.compress_prog, @@ -631,3 +651,68 @@ impl MergeInput for PlainMergeInput { &mut self.inner } } + +#[cfg(all(test, target_os = "linux"))] +mod tests { + use std::fs::{self, File}; + use std::io::Write; + use std::os::unix::fs::PermissionsExt; + use std::path::PathBuf; + use std::thread; + use std::time::{Duration, Instant}; + + use super::{WriteableCompressedTmpFile, WriteablePlainTmpFile, WriteableTmpFile}; + + #[test] + fn plain_tmp_file_propagates_flush_errors() { + let file = File::options().write(true).open("/dev/full").unwrap(); + let mut tmp_file = + WriteablePlainTmpFile::create((file, PathBuf::from("/dev/full")), None).unwrap(); + tmp_file.as_write().write_all(b"buffered data").unwrap(); + + let error = tmp_file.finished_writing().err().unwrap(); + + assert_eq!(error.to_string(), "No space left on device"); + } + + #[test] + fn compressed_tmp_file_propagates_flush_errors() { + let directory = tempfile::tempdir().unwrap(); + let compressor = directory.path().join("compressor.sh"); + let ready = directory.path().join("ready"); + let done = directory.path().join("done"); + fs::write( + &compressor, + format!( + "#!/bin/sh\nexec 0<&-\n: > '{}'\nsleep 1\n: > '{}'\nexit 1\n", + ready.display(), + done.display() + ), + ) + .unwrap(); + let mut permissions = fs::metadata(&compressor).unwrap().permissions(); + permissions.set_mode(0o755); + fs::set_permissions(&compressor, permissions).unwrap(); + + let output = tempfile::tempfile().unwrap(); + let mut tmp_file = WriteableCompressedTmpFile::create( + (output, PathBuf::new()), + Some(compressor.to_str().unwrap()), + ) + .unwrap(); + let child_pid = tmp_file.child.id(); + tmp_file.as_write().write_all(b"buffered data").unwrap(); + + let deadline = Instant::now() + Duration::from_secs(5); + while !ready.exists() { + assert!(Instant::now() < deadline, "compressor did not become ready"); + thread::sleep(Duration::from_millis(10)); + } + + let error = tmp_file.finished_writing().err().unwrap(); + + assert_eq!(error.to_string(), "Broken pipe"); + assert!(done.exists()); + assert!(!PathBuf::from(format!("/proc/{child_pid}")).exists()); + } +} diff --git a/tests/by-util/test_sort.rs b/tests/by-util/test_sort.rs index cf0244889cd..11fa48fe138 100644 --- a/tests/by-util/test_sort.rs +++ b/tests/by-util/test_sort.rs @@ -1182,6 +1182,22 @@ fn test_merge_write_error_does_not_panic() { } } +#[test] +#[cfg(target_os = "linux")] +fn test_merge_flush_error_is_reported() { + use std::fs::File; + + let ts = TestScenario::new("sort"); + ts.fixtures.write("input.txt", "line\n"); + + let dev_full = File::create("/dev/full").expect("Failed to open /dev/full"); + ts.ucmd() + .args(&["-m", "input.txt"]) + .set_stdout(dev_full) + .fails() + .stderr_contains("No space left on device"); +} + #[test] fn test_merge_unique() { new_ucmd!() From 4265a9cb5de29e35449669dbe85dc4c738bb3320 Mon Sep 17 00:00:00 2001 From: Anthony DePasquale Date: Thu, 13 Aug 2026 11:22:57 +0200 Subject: [PATCH 2/2] sort: limit fix to final merge flush --- src/uu/sort/src/ext_sort/threaded.rs | 53 ++-------------- src/uu/sort/src/merge.rs | 93 +++------------------------- 2 files changed, 13 insertions(+), 133 deletions(-) diff --git a/src/uu/sort/src/ext_sort/threaded.rs b/src/uu/sort/src/ext_sort/threaded.rs index c4ef133d1f0..7dd089d0fe8 100644 --- a/src/uu/sort/src/ext_sort/threaded.rs +++ b/src/uu/sort/src/ext_sort/threaded.rs @@ -293,56 +293,13 @@ fn write( separator: u8, ) -> UResult { let mut tmp_file = I::create(file, compress_prog)?; - let write_result = write_lines(chunk.lines(), tmp_file.as_write(), separator); - let finish_result = tmp_file.finished_writing(); - write_result?; - finish_result + write_lines(chunk.lines(), tmp_file.as_write(), separator); + tmp_file.finished_writing() } -fn write_lines(lines: &[Line], writer: &mut T, separator: u8) -> std::io::Result<()> { +fn write_lines(lines: &[Line], writer: &mut T, separator: u8) { for s in lines { - writer.write_all(s.line)?; - writer.write_all(&[separator])?; - } - Ok(()) -} - -#[cfg(test)] -mod tests { - use std::io::{self, Write}; - - use super::write_lines; - use crate::Line; - - struct FailAfterFirstWrite { - writes: usize, - } - - impl Write for FailAfterFirstWrite { - fn write(&mut self, buf: &[u8]) -> io::Result { - if self.writes == 0 { - self.writes += 1; - Ok(buf.len()) - } else { - Err(io::Error::other("write failed")) - } - } - - fn flush(&mut self) -> io::Result<()> { - Ok(()) - } - } - - #[test] - fn write_lines_propagates_write_errors() { - let lines = [Line { - line: b"line", - index: 0, - }]; - let mut writer = FailAfterFirstWrite { writes: 0 }; - - let error = write_lines(&lines, &mut writer, b'\n').unwrap_err(); - - assert_eq!(error.kind(), io::ErrorKind::Other); + writer.write_all(s.line).unwrap(); + writer.write_all(&[separator]).unwrap(); } } diff --git a/src/uu/sort/src/merge.rs b/src/uu/sort/src/merge.rs index bf85e49fc5f..b527e5fa877 100644 --- a/src/uu/sort/src/merge.rs +++ b/src/uu/sort/src/merge.rs @@ -142,10 +142,8 @@ pub fn merge_with_file_limit< let mut tmp_file = Tmp::create(tmp_dir.next_file()?, settings.compress_prog.as_deref())?; - let write_result = merger.write_all_to(settings, tmp_file.as_write()); - let finish_result = tmp_file.finished_writing(); - write_result?; - temporary_files.push(finish_result?); + merger.write_all_to(settings, tmp_file.as_write())?; + temporary_files.push(tmp_file.finished_writing()?); } } // Merge any remaining files that didn't get merged in a full batch above. @@ -155,10 +153,8 @@ pub fn merge_with_file_limit< let mut tmp_file = Tmp::create(tmp_dir.next_file()?, settings.compress_prog.as_deref())?; - let write_result = merger.write_all_to(settings, tmp_file.as_write()); - let finish_result = tmp_file.finished_writing(); - write_result?; - temporary_files.push(finish_result?); + merger.write_all_to(settings, tmp_file.as_write())?; + temporary_files.push(tmp_file.finished_writing()?); } merge_with_file_limit::<_, _, Tmp>( temporary_files @@ -317,7 +313,7 @@ impl FileMerger<'_> { let ctx = || translate!("sort-error-write-failed", "output" => output_name.maybe_quote()); let mut out = output.into_write(); self.write_all_to(settings, &mut out)?; - flush_writer(&mut out).map_err_context(ctx) + out.flush().map_err_context(ctx) } fn write_all_to(mut self, settings: &GlobalSettings, out: &mut impl Write) -> UResult<()> { @@ -459,10 +455,6 @@ fn check_child_success(mut child: Child, program: &str) -> UResult<()> { } } -fn flush_writer(writer: &mut impl Write) -> std::io::Result<()> { - writer.flush() -} - /// A temporary file that can be written to. pub trait WriteableTmpFile: Sized { type Closed: ClosedTmpFile; @@ -509,8 +501,7 @@ impl WriteableTmpFile for WriteablePlainTmpFile { }) } - fn finished_writing(mut self) -> UResult { - flush_writer(&mut self.file)?; + fn finished_writing(self) -> UResult { Ok(ClosedPlainTmpFile { path: self.path }) } @@ -582,12 +573,9 @@ impl WriteableTmpFile for WriteableCompressedTmpFile { }) } - fn finished_writing(mut self) -> UResult { - let flush_result = flush_writer(&mut self.child_stdin); + fn finished_writing(self) -> UResult { drop(self.child_stdin); - let child_result = check_child_success(self.child, &self.compress_prog); - flush_result?; - child_result?; + check_child_success(self.child, &self.compress_prog)?; Ok(ClosedCompressedTmpFile { path: self.path, compress_prog: self.compress_prog, @@ -651,68 +639,3 @@ impl MergeInput for PlainMergeInput { &mut self.inner } } - -#[cfg(all(test, target_os = "linux"))] -mod tests { - use std::fs::{self, File}; - use std::io::Write; - use std::os::unix::fs::PermissionsExt; - use std::path::PathBuf; - use std::thread; - use std::time::{Duration, Instant}; - - use super::{WriteableCompressedTmpFile, WriteablePlainTmpFile, WriteableTmpFile}; - - #[test] - fn plain_tmp_file_propagates_flush_errors() { - let file = File::options().write(true).open("/dev/full").unwrap(); - let mut tmp_file = - WriteablePlainTmpFile::create((file, PathBuf::from("/dev/full")), None).unwrap(); - tmp_file.as_write().write_all(b"buffered data").unwrap(); - - let error = tmp_file.finished_writing().err().unwrap(); - - assert_eq!(error.to_string(), "No space left on device"); - } - - #[test] - fn compressed_tmp_file_propagates_flush_errors() { - let directory = tempfile::tempdir().unwrap(); - let compressor = directory.path().join("compressor.sh"); - let ready = directory.path().join("ready"); - let done = directory.path().join("done"); - fs::write( - &compressor, - format!( - "#!/bin/sh\nexec 0<&-\n: > '{}'\nsleep 1\n: > '{}'\nexit 1\n", - ready.display(), - done.display() - ), - ) - .unwrap(); - let mut permissions = fs::metadata(&compressor).unwrap().permissions(); - permissions.set_mode(0o755); - fs::set_permissions(&compressor, permissions).unwrap(); - - let output = tempfile::tempfile().unwrap(); - let mut tmp_file = WriteableCompressedTmpFile::create( - (output, PathBuf::new()), - Some(compressor.to_str().unwrap()), - ) - .unwrap(); - let child_pid = tmp_file.child.id(); - tmp_file.as_write().write_all(b"buffered data").unwrap(); - - let deadline = Instant::now() + Duration::from_secs(5); - while !ready.exists() { - assert!(Instant::now() < deadline, "compressor did not become ready"); - thread::sleep(Duration::from_millis(10)); - } - - let error = tmp_file.finished_writing().err().unwrap(); - - assert_eq!(error.to_string(), "Broken pipe"); - assert!(done.exists()); - assert!(!PathBuf::from(format!("/proc/{child_pid}")).exists()); - } -}