From 13b1b5b134fe4abdfea6f2f6cc2da6490d2c41e4 Mon Sep 17 00:00:00 2001 From: can1357 Date: Tue, 30 Jun 2026 23:45:42 +0200 Subject: [PATCH] refactor: migrated concurrency primitives to flume and parking_lot - Replaced standard and tokio mpsc channels with flume channels across workspace crates to simplify thread synchronization. - Swapped standard Mutex guards for parking_lot Mutexes to avoid manual lock poisoning handling and improve performance. - Declared workspace-wide dependencies for flume and parking_lot in root and member Cargo manifests. --- Cargo.lock | 38 ++++++++ Cargo.toml | 1 + crates/pi-natives/Cargo.toml | 1 + crates/pi-natives/src/appearance.rs | 4 +- crates/pi-natives/src/clipboard.rs | 6 +- crates/pi-natives/src/grep.rs | 38 +++----- crates/pi-natives/src/pty.rs | 43 ++++----- crates/pi-natives/src/shell.rs | 15 ++- crates/pi-shell/Cargo.toml | 2 + crates/pi-shell/src/coreutils.rs | 8 +- crates/pi-shell/src/process.rs | 15 +-- crates/pi-shell/src/shell.rs | 94 +++++++++---------- crates/pi-uu-grep/Cargo.toml | 3 + crates/pi-uu-grep/src/lib.rs | 11 +-- crates/pi-walker/Cargo.toml | 1 + crates/pi-walker/src/cache.rs | 22 ++--- crates/pi-walker/src/lib.rs | 14 +-- crates/vendor/uu-sort/Cargo.toml | 1 + crates/vendor/uu-sort/src/check.rs | 20 ++-- crates/vendor/uu-sort/src/chunks.rs | 4 +- .../vendor/uu-sort/src/ext_sort/threaded.rs | 12 +-- crates/vendor/uu-sort/src/merge.rs | 10 +- 22 files changed, 180 insertions(+), 183 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index d74bcdd5e..5a08f6926 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1213,6 +1213,18 @@ dependencies = [ "thiserror 2.0.18", ] +[[package]] +name = "flume" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "da0e4dd2a88388a1f4ccc7c9ce104604dab68d9f408dc34cd45823d5a9069095" +dependencies = [ + "futures-core", + "futures-sink", + "nanorand", + "spin", +] + [[package]] name = "fnv" version = "1.0.7" @@ -1382,8 +1394,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff2abc00be7fca6ebc474524697ae276ad847ad0a6b3faa4bcb027e9a4614ad0" dependencies = [ "cfg-if", + "js-sys", "libc", "wasi 0.11.1+wasi-snapshot-preview1", + "wasm-bindgen", ] [[package]] @@ -2299,6 +2313,15 @@ dependencies = [ "pxfm", ] +[[package]] +name = "nanorand" +version = "0.7.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6a51313c5820b0b02bd422f4b44776fbf47961755c74ce64afc73bfad10226c3" +dependencies = [ + "getrandom 0.2.17", +] + [[package]] name = "napi" version = "3.9.4" @@ -3000,6 +3023,7 @@ dependencies = [ "ast-grep-core", "base64", "clap", + "flume", "fontdue", "globset", "grep-matcher", @@ -3052,10 +3076,12 @@ dependencies = [ "brush-parser 0.3.0", "bytes", "clap", + "flume", "globset", "ignore", "libc", "os_pipe", + "parking_lot", "pi-uutils-ctx", "pi-walker", "pi_uu_grep", @@ -3097,6 +3123,7 @@ dependencies = [ "globset", "ignore", "libc", + "parking_lot", "rayon", "windows-sys 0.61.2", ] @@ -3111,6 +3138,7 @@ dependencies = [ "grep-regex", "grep-searcher", "ignore", + "parking_lot", "pi-uutils-ctx", "pi-walker", ] @@ -3767,6 +3795,15 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "spin" +version = "0.9.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67" +dependencies = [ + "lock_api", +] + [[package]] name = "stable_deref_trait" version = "1.2.1" @@ -4966,6 +5003,7 @@ dependencies = [ "binary-heap-plus", "clap", "compare", + "flume", "foldhash 0.2.0", "itertools", "memchr", diff --git a/Cargo.toml b/Cargo.toml index f4bfbf11c..d6a69e456 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -184,6 +184,7 @@ parking_lot = "0.12.5" rayon = "1.12" tokio = { version = "1", features = ["full"] } tokio-util = { version = "0.7", features = ["full"] } +flume = "0.11" # ────────────────────────────────────────────────────────────────────────────── # Serialization & Data Formats diff --git a/crates/pi-natives/Cargo.toml b/crates/pi-natives/Cargo.toml index 0928d3dc3..6495a02c8 100644 --- a/crates/pi-natives/Cargo.toml +++ b/crates/pi-natives/Cargo.toml @@ -33,6 +33,7 @@ napi.workspace = true napi-derive.workspace = true parking_lot.workspace = true phf.workspace = true +flume.workspace = true pi-ast.workspace = true pi-iso.workspace = true pi-shell.workspace = true diff --git a/crates/pi-natives/src/appearance.rs b/crates/pi-natives/src/appearance.rs index bae49083f..e27e855b8 100644 --- a/crates/pi-natives/src/appearance.rs +++ b/crates/pi-natives/src/appearance.rs @@ -35,7 +35,7 @@ mod platform { use std::{ ffi::{CStr, CString, c_char, c_void}, ptr, - sync::{Arc, mpsc}, + sync::Arc, thread::{self, JoinHandle}, }; @@ -285,7 +285,7 @@ mod platform { let rl_clone = run_loop.clone(); // Signal that the background thread has stored its `CFRunLoopRef`. - let (tx, rx) = mpsc::sync_channel::<()>(1); + let (tx, rx) = flume::bounded::<()>(1); let handle = thread::spawn(move || { // SAFETY: All CoreFoundation objects created or copied here are either released diff --git a/crates/pi-natives/src/clipboard.rs b/crates/pi-natives/src/clipboard.rs index 792e36d50..060cae6c7 100644 --- a/crates/pi-natives/src/clipboard.rs +++ b/crates/pi-natives/src/clipboard.rs @@ -65,11 +65,13 @@ pub fn copy_to_clipboard(text: String) -> Result<()> { /// is harmless there. #[cfg(target_os = "linux")] fn set_clipboard_text(text: String) -> Result<()> { - use std::sync::{Mutex, OnceLock}; + use std::sync::OnceLock; + + use parking_lot::Mutex; static CLIPBOARD: OnceLock>> = OnceLock::new(); let cell = CLIPBOARD.get_or_init(|| Mutex::new(None)); - let mut guard = cell.lock().unwrap_or_else(|poisoned| poisoned.into_inner()); + let mut guard = cell.lock(); if guard.is_none() { *guard = Some( Clipboard::new() diff --git a/crates/pi-natives/src/grep.rs b/crates/pi-natives/src/grep.rs index fb3d65f8c..ec61534b6 100644 --- a/crates/pi-natives/src/grep.rs +++ b/crates/pi-natives/src/grep.rs @@ -12,10 +12,7 @@ use std::{ fs::File, io::{self, Read}, path::{Path, PathBuf}, - sync::{ - Mutex, - atomic::{AtomicU64, AtomicUsize, Ordering}, - }, + sync::atomic::{AtomicU64, AtomicUsize, Ordering}, }; use grep_matcher::Matcher; @@ -29,6 +26,7 @@ use napi::{ threadsafe_function::{ThreadsafeFunction, ThreadsafeFunctionCallMode}, }; use napi_derive::napi; +use parking_lot::Mutex; use smallvec::SmallVec; use crate::{glob_util, iofs, task}; @@ -1234,11 +1232,7 @@ fn handle_file( } match search_one_file(searcher, matcher, file, file_params, policy) { FileOutcome::Defer => { - state - .deferred - .lock() - .expect("deferred lock poisoned") - .push(file.clone()); + state.deferred.lock().push(file.clone()); }, FileOutcome::SkippedOversized => { state.skipped_oversized.fetch_add(1, Ordering::Relaxed); @@ -1248,16 +1242,12 @@ fn handle_file( state.files_searched.fetch_add(1, Ordering::Relaxed); if search.match_count > 0 { let emitted_in_file = search.collected; - state - .results - .lock() - .expect("results lock poisoned") - .push(FileSearchResult { - relative_path: file.relative.clone(), - matches: search.matches, - match_count: search.match_count, - limit_reached: search.limit_reached, - }); + state.results.lock().push(FileSearchResult { + relative_path: file.relative.clone(), + matches: search.matches, + match_count: search.match_count, + limit_reached: search.limit_reached, + }); if stop_after_matches.is_some() { state.emitted.fetch_add(emitted_in_file, Ordering::Relaxed); } @@ -1316,7 +1306,7 @@ fn run_pass( )?; } } - let mut results = std::mem::take(&mut *state.results.lock().expect("results lock poisoned")); + let mut results = std::mem::take(&mut *state.results.lock()); results.sort_unstable_by(|a, b| a.relative_path.cmp(&b.relative_path)); Ok(results) } @@ -1348,11 +1338,7 @@ fn process_candidates( None => true, }); if !oversized_hinted.is_empty() { - state - .deferred - .lock() - .expect("deferred lock poisoned") - .extend(oversized_hinted); + state.deferred.lock().extend(oversized_hinted); } let mut results = run_pass( @@ -1368,7 +1354,7 @@ fn process_candidates( // Pass 2: deferred oversized files, searched over their leading window — // only when a content-mode budget was not already satisfied in pass 1. - let deferred = std::mem::take(&mut *state.deferred.lock().expect("deferred lock poisoned")); + let deferred = std::mem::take(&mut *state.deferred.lock()); let limit_satisfied = stop_after_matches.is_some_and(|stop| state.emitted.load(Ordering::Relaxed) >= stop); if !deferred.is_empty() && !limit_satisfied { diff --git a/crates/pi-natives/src/pty.rs b/crates/pi-natives/src/pty.rs index 391100b57..c0354d03f 100644 --- a/crates/pi-natives/src/pty.rs +++ b/crates/pi-natives/src/pty.rs @@ -8,7 +8,7 @@ use std::{ collections::HashMap, io::{Read, Write}, str, - sync::{Arc, Mutex, mpsc}, + sync::Arc, time::{Duration, Instant}, }; @@ -17,6 +17,7 @@ use napi::{ threadsafe_function::{ThreadsafeFunction, ThreadsafeFunctionCallMode}, }; use napi_derive::napi; +use parking_lot::Mutex; use portable_pty::{Child, CommandBuilder, PtySize, native_pty_system}; use crate::{ps, task}; @@ -83,7 +84,7 @@ const POST_EXIT_DRAIN_TIMEOUT: Duration = Duration::from_millis(300); const FINAL_READER_DRAIN_TIMEOUT: Duration = Duration::from_millis(50); struct PtySessionCore { - control_tx: mpsc::Sender, + control_tx: flume::Sender, } /// Stateful PTY session for interactive stdin/stdout passthrough. @@ -126,11 +127,9 @@ impl PtySession { let core = Arc::clone(&self.core); // Register control channel synchronously so write()/kill() work immediately. - let (control_tx, control_rx) = mpsc::channel::(); + let (control_tx, control_rx) = flume::unbounded::(); { - let mut guard = core - .lock() - .map_err(|_| Error::from_reason("PTY session lock poisoned"))?; + let mut guard = core.lock(); if guard.is_some() { return Err(Error::from_reason("PTY session already running")); } @@ -142,9 +141,7 @@ impl PtySession { .await; // Always clear core regardless of result - let mut guard = core - .lock() - .map_err(|_| Error::from_reason("PTY session lock poisoned"))?; + let mut guard = core.lock(); *guard = None; drop(guard); @@ -179,10 +176,7 @@ impl PtySession { impl PtySession { fn send_control(&self, message: ControlMessage) -> Result<()> { - let guard = self - .core - .lock() - .map_err(|_| Error::from_reason("PTY session lock poisoned"))?; + let guard = self.core.lock(); let core = guard .as_ref() .ok_or_else(|| Error::from_reason("PTY session is not running"))?; @@ -213,7 +207,7 @@ fn terminate_pty_processes( fn run_pty_sync( config: PtyRunConfig, on_chunk: Option>, - control_rx: mpsc::Receiver, + control_rx: flume::Receiver, ct: task::CancelToken, ) -> Result { let pty_system = native_pty_system(); @@ -225,7 +219,7 @@ fn run_pty_sync( // Windows ConPTY openpty() can hang indefinitely when the console // subsystem isn't properly initialized. Use a short startup timeout // so the Promise rejects instead of hanging forever. - let (tx, rx) = mpsc::channel(); + let (tx, rx) = flume::unbounded(); std::thread::spawn(move || { let result = pty_system.openpty(PtySize { rows: config.rows, @@ -303,7 +297,7 @@ fn run_pty_sync( .try_clone_reader() .map_err(|err| Error::from_reason(format!("Failed to create PTY reader: {err}")))?; - let (reader_tx, reader_rx) = mpsc::channel::(); + let (reader_tx, reader_rx) = flume::unbounded::(); let reader_thread = std::thread::spawn(move || { const REPLACEMENT: &str = "\u{FFFD}"; const BUF: usize = 65536; @@ -404,8 +398,7 @@ fn run_pty_sync( reader_drain_deadline = Some(Instant::now() + POST_CANCEL_DRAIN_TIMEOUT); } }, - Err(mpsc::TryRecvError::Empty) => break, - Err(mpsc::TryRecvError::Disconnected) => break, + Err(flume::TryRecvError::Empty | flume::TryRecvError::Disconnected) => break, } } @@ -416,8 +409,8 @@ fn run_pty_sync( reader_done = true; break; }, - Err(mpsc::TryRecvError::Empty) => break, - Err(mpsc::TryRecvError::Disconnected) => { + Err(flume::TryRecvError::Empty) => break, + Err(flume::TryRecvError::Disconnected) => { reader_done = true; break; }, @@ -448,8 +441,8 @@ fn run_pty_sync( match reader_rx.recv_timeout(wait_duration) { Ok(ReaderEvent::Chunk(chunk)) => emit_chunk(&chunk, on_chunk.as_ref()), Ok(ReaderEvent::Done) => reader_done = true, - Err(mpsc::RecvTimeoutError::Timeout) => {}, - Err(mpsc::RecvTimeoutError::Disconnected) => { + Err(flume::RecvTimeoutError::Timeout) => {}, + Err(flume::RecvTimeoutError::Disconnected) => { reader_done = true; if exit_code.is_none() { std::thread::sleep(wait_duration); @@ -519,8 +512,8 @@ fn run_pty_sync( reader_done = true; break; }, - Err(mpsc::RecvTimeoutError::Timeout) => {}, - Err(mpsc::RecvTimeoutError::Disconnected) => { + Err(flume::RecvTimeoutError::Timeout) => {}, + Err(flume::RecvTimeoutError::Disconnected) => { reader_done = true; break; }, @@ -537,7 +530,7 @@ fn run_pty_sync( // but the main thread never blocks. #[cfg(windows)] { - let (drop_tx, drop_rx) = mpsc::channel::<()>(); + let (drop_tx, drop_rx) = flume::unbounded::<()>(); std::thread::spawn(move || { drop(master); let _ = drop_tx.send(()); diff --git a/crates/pi-natives/src/shell.rs b/crates/pi-natives/src/shell.rs index 3958fac07..2b55d56d8 100644 --- a/crates/pi-natives/src/shell.rs +++ b/crates/pi-natives/src/shell.rs @@ -6,7 +6,6 @@ use napi::{ Env, Result, bindgen_prelude::*, threadsafe_function::{ThreadsafeFunction, ThreadsafeFunctionCallMode}, - tokio::sync::mpsc, }; use napi_derive::napi; use pi_shell::{ @@ -293,11 +292,11 @@ pub fn execute_shell<'env>( fn bridge_chunks( on_chunk: Option>, -) -> (Option>, Option>) { +) -> (Option>, Option>) { let Some(on_chunk) = on_chunk else { return (None, None); }; - let (tx, mut rx) = mpsc::unbounded_channel::(); + let (tx, rx) = flume::unbounded::(); let handle = napi::tokio::spawn(async move { // Hard cap on one coalesced batch so the JS main thread never sees a // multi-MB napi callback (a giant single string would stall sanitize + @@ -307,7 +306,7 @@ fn bridge_chunks( // each batch because `String` ownership is moved into the napi call. const INITIAL_BATCH_CAP: usize = 8 * 1024; let mut batch = String::with_capacity(INITIAL_BATCH_CAP); - while let Some(first) = rx.recv().await { + while let Ok(first) = rx.recv_async().await { batch.push_str(&first); // Greedily drain everything already queued. Child processes that // write byte-at-a-time (printf-style progress, llama-cli token @@ -357,12 +356,12 @@ pub fn apply_bash_fixups(command: String) -> BashFixupResult { mod tests { use std::time::Duration; + #[cfg(unix)] + use flume; use pi_shell::{ ShellRunOptions as CoreShellRunOptions, cancel::{AbortReason, CancelToken}, }; - #[cfg(unix)] - use tokio::sync::mpsc; use tokio::time; use super::CoreShell; @@ -412,7 +411,7 @@ mod tests { #[tokio::test(flavor = "multi_thread")] async fn embedded_external_command_runs_in_its_own_session() { let shell = CoreShell::new(None); - let (tx, mut rx) = mpsc::unbounded_channel::(); + let (tx, rx) = flume::unbounded::(); let handle = tokio::spawn(async move { shell .run( @@ -427,7 +426,7 @@ mod tests { ) .await }); - let child_pid = time::timeout(Duration::from_secs(5), rx.recv()) + let child_pid = time::timeout(Duration::from_secs(5), rx.recv_async()) .await .expect("timed out waiting for child pid") .expect("missing child pid chunk") diff --git a/crates/pi-shell/Cargo.toml b/crates/pi-shell/Cargo.toml index a8f39fe9d..9ed381ca7 100644 --- a/crates/pi-shell/Cargo.toml +++ b/crates/pi-shell/Cargo.toml @@ -21,6 +21,7 @@ clap.workspace = true os_pipe.workspace = true globset.workspace = true ignore.workspace = true +parking_lot.workspace = true regex.workspace = true serde.workspace = true serde_json.workspace = true @@ -28,6 +29,7 @@ tokio.workspace = true tokio-util.workspace = true toml.workspace = true xxhash-rust.workspace = true +flume.workspace = true pi-walker.workspace = true pi-uutils-ctx = { path = "../pi-uutils-ctx" } uu_mkdir = { path = "../vendor/uu-mkdir" } diff --git a/crates/pi-shell/src/coreutils.rs b/crates/pi-shell/src/coreutils.rs index 7bd6fca9c..10f1c56a2 100644 --- a/crates/pi-shell/src/coreutils.rs +++ b/crates/pi-shell/src/coreutils.rs @@ -217,14 +217,16 @@ mod tests { ffi::OsString, io::{self, Write}, path::PathBuf, - sync::{Arc, atomic::AtomicBool, mpsc}, + sync::{Arc, atomic::AtomicBool}, }; + use flume::Sender; + use super::{UutilRun, run_caught}; /// `Send` writer that forwards every write onto a channel so a test can /// inspect what the utility wrote to the scope's stderr. - struct ChanWriter(mpsc::Sender>); + struct ChanWriter(Sender>); impl Write for ChanWriter { fn write(&mut self, buf: &[u8]) -> io::Result { let _ = self.0.send(buf.to_vec()); @@ -250,7 +252,7 @@ mod tests { } fn run_in_scope(run: UutilRun, argv: Vec) -> (i32, String) { - let (tx, rx) = mpsc::channel(); + let (tx, rx) = flume::unbounded(); let code = pi_uutils_ctx::scope(scope_io(Box::new(ChanWriter(tx))), || run_caught(run, argv)); let mut err = Vec::new(); while let Ok(chunk) = rx.try_recv() { diff --git a/crates/pi-shell/src/process.rs b/crates/pi-shell/src/process.rs index 48a24a6f3..ede0fd3f8 100644 --- a/crates/pi-shell/src/process.rs +++ b/crates/pi-shell/src/process.rs @@ -3,6 +3,7 @@ use std::{collections::HashSet, time::Duration}; use anyhow::Result; +use parking_lot::Mutex; use crate::cancel::CancelToken; @@ -1614,7 +1615,7 @@ struct SpawnedProcess { /// explicit — only processes this run actually spawned are ever signalled. #[derive(Default)] pub struct SpawnRegistry { - spawned: std::sync::Mutex>, + spawned: Mutex>, } impl SpawnRegistry { @@ -1626,11 +1627,7 @@ impl SpawnRegistry { /// Record a freshly spawned child. Called from the spawn-observer hook. pub fn record(&self, pid: i32, pgid: Option) { - self - .spawned - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .push(SpawnedProcess { pid, pgid }); + self.spawned.lock().push(SpawnedProcess { pid, pgid }); } /// Build the kill set from the processes recorded so far. Re-read on every @@ -1644,11 +1641,7 @@ impl SpawnRegistry { #[must_use] pub fn build_targets(&self) -> TerminationTargets { let mut targets = TerminationTargets::new(); - let spawned = self - .spawned - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .clone(); + let spawned = self.spawned.lock().clone(); for entry in spawned { targets.add_pid(entry.pid); if let Some(pgid) = entry.pgid diff --git a/crates/pi-shell/src/shell.rs b/crates/pi-shell/src/shell.rs index 4b02f3dde..608d9da04 100644 --- a/crates/pi-shell/src/shell.rs +++ b/crates/pi-shell/src/shell.rs @@ -22,12 +22,10 @@ use brush_core::{ }; use bytes::Bytes; use clap::Parser; +use flume::Sender; #[cfg(not(unix))] use tokio::io::AsyncReadExt as _; -use tokio::{ - sync::{Mutex as TokioMutex, mpsc}, - time, -}; +use tokio::{sync::Mutex as TokioMutex, time}; use tokio_util::sync::CancellationToken; #[cfg(windows)] @@ -153,7 +151,7 @@ impl Shell { pub async fn run( &self, options: ShellRunOptions, - on_chunk: Option>, + on_chunk: Option>, mut cancel_token: CancelToken, ) -> Result { let run_config = ShellRunConfig { @@ -209,7 +207,7 @@ impl Shell { pub async fn execute_shell( options: ShellExecuteOptions, - on_chunk: Option>, + on_chunk: Option>, cancel_token: CancelToken, ) -> Result { let minimizer = options @@ -234,8 +232,8 @@ pub async fn execute_shell( /// its bytes are dropped. #[derive(Default)] pub struct StreamSinks { - pub stdout: Option>, - pub stderr: Option>, + pub stdout: Option>, + pub stderr: Option>, } /// One-shot execution that delivers stdout/stderr as raw byte chunks. @@ -267,7 +265,7 @@ async fn run_shell_session( abort_state: ShellAbortState, config: ShellConfig, run_config: ShellRunConfig, - on_chunk: Option>, + on_chunk: Option>, ct: &mut CancelToken, ) -> Result { let tokio_cancel = CancellationToken::new(); @@ -353,7 +351,7 @@ async fn run_shell_session( async fn run_shell_oneshot( config: ShellConfig, run_config: ShellRunConfig, - on_chunk: Option>, + on_chunk: Option>, ct: CancelToken, ) -> Result { let tokio_cancel = CancellationToken::new(); @@ -759,7 +757,7 @@ impl ChainCapture { async fn run_shell_command( session: &mut ShellSessionCore, options: &ShellRunConfig, - on_chunk: Option>, + on_chunk: Option>, cancel_token: CancellationToken, spawn_registry: Arc, ) -> Result<(ExecutionResult, Option)> { @@ -810,7 +808,7 @@ async fn run_shell_command( async fn run_shell_command_single( session: &mut ShellSessionCore, options: &ShellRunConfig, - on_chunk: Option>, + on_chunk: Option>, cancel_token: CancellationToken, spawn_registry: Arc, minimizer_mode: minimizer::engine::MinimizerMode, @@ -893,7 +891,7 @@ async fn run_shell_command_single( async fn run_shell_command_segmented_chain( session: &mut ShellSessionCore, options: &ShellRunConfig, - on_chunk: Option>, + on_chunk: Option>, cancel_token: CancellationToken, spawn_registry: Arc, ) -> Result<(ExecutionResult, Option)> { @@ -1035,7 +1033,7 @@ async fn run_shell_command_once( session: &mut ShellSessionCore, mut command: String, mut params: ExecutionParameters, - on_chunk: Option>, + on_chunk: Option>, cancel_token: CancellationToken, spawn_registry: Arc, capture_mode: CommandCaptureMode, @@ -1056,7 +1054,7 @@ async fn run_shell_command_once( params.set_cancel_token(cancel_token.clone()); params.set_spawn_observer(spawn_registry.clone()); let reader_cancel = CancellationToken::new(); - let (activity_tx, mut activity_rx) = mpsc::channel::<()>(1); + let (activity_tx, activity_rx) = flume::bounded::<()>(1); let reader_callback = on_chunk; let mut reader_handle = tokio::spawn({ let reader_cancel = reader_cancel.clone(); @@ -1123,8 +1121,8 @@ async fn run_shell_command_once( reader_finished = true; break; } - msg = activity_rx.recv() => { - if msg.is_none() { + msg = activity_rx.recv_async() => { + if msg.is_err() { break; } idle_timer.as_mut().reset(time::Instant::now() + POST_EXIT_IDLE); @@ -1188,7 +1186,7 @@ async fn run_shell_command_streams( params.set_cancel_token(cancel_token.clone()); params.set_spawn_observer(spawn_registry.clone()); let reader_cancel = CancellationToken::new(); - let (activity_tx, mut activity_rx) = mpsc::channel::<()>(1); + let (activity_tx, activity_rx) = flume::bounded::<()>(1); let StreamSinks { stdout: stdout_sink, stderr: stderr_sink } = streams; let mut stdout_handle = tokio::spawn(Box::pin(read_output_bytes( @@ -1256,8 +1254,8 @@ async fn run_shell_command_streams( let _ = res; stderr_finished = true; } - msg = activity_rx.recv() => { - if msg.is_none() { + msg = activity_rx.recv_async() => { + if msg.is_err() { break; } idle_timer.as_mut().reset(time::Instant::now() + POST_EXIT_IDLE); @@ -1295,9 +1293,9 @@ async fn run_shell_command_streams( async fn read_output_bytes( reader: fs::File, - sink: Option>, + sink: Option>, cancel_token: CancellationToken, - activity: mpsc::Sender<()>, + activity: Sender<()>, ) { const BUF: usize = 65536; @@ -1564,9 +1562,9 @@ struct BufferedOutput { async fn read_output( reader: fs::File, - on_chunk: Option>, + on_chunk: Option>, cancel_token: CancellationToken, - activity: mpsc::Sender<()>, + activity: Sender<()>, ) { const REPLACEMENT: &str = "\u{FFFD}"; const BUF: usize = 65536; @@ -1672,9 +1670,9 @@ async fn read_output( async fn read_output_buffered( reader: fs::File, - on_chunk: Option>, + on_chunk: Option>, cancel_token: CancellationToken, - activity: mpsc::Sender<()>, + activity: Sender<()>, max_capture_bytes: usize, ) -> BufferedOutput { const REPLACEMENT: &str = "\u{FFFD}"; @@ -1832,7 +1830,7 @@ fn read_nonblocking(file: &T, buf: &mut [u8]) -> io::Re } } -fn emit_chunk(text: &str, callback: Option<&mpsc::UnboundedSender>) { +fn emit_chunk(text: &str, callback: Option<&Sender>) { if let Some(callback) = callback { let _ = callback.send(text.to_string()); } @@ -2980,7 +2978,7 @@ mod tests { cancel_token: CancelToken, ) -> (ShellExecuteResult, String) { let _guard = shell_test_lock().lock().await; - let (tx, mut rx) = mpsc::unbounded_channel::(); + let (tx, rx) = flume::unbounded::(); let options = ShellExecuteOptions { command: command.to_string(), cwd: cwd.map(|path| path.to_string_lossy().into_owned()), @@ -2991,7 +2989,7 @@ mod tests { .await .expect("execute_shell"); let mut output = String::new(); - while let Some(chunk) = rx.recv().await { + while let Ok(chunk) = rx.recv_async().await { output.push_str(&chunk); } (result, output) @@ -3476,7 +3474,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] let _guard = shell_test_lock().lock().await; let shell_b = Shell::new(None); - let (tx_b, mut rx_b) = mpsc::unbounded_channel::(); + let (tx_b, rx_b) = flume::unbounded::(); let mut ct_b = CancelToken::default(); let abort_b = ct_b.emplace_abort_token(); let handle_b = tokio::spawn(async move { @@ -3496,7 +3494,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] let b_ready = time::timeout(Duration::from_secs(5), async { loop { let chunk = rx_b - .recv() + .recv_async() .await .expect("run B ended before printing readiness"); b_output.push_str(&chunk); @@ -3510,7 +3508,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] assert_eq!(b_ready.trim(), "ready", "run B should reach its long sleep before run A starts"); let shell_a = Shell::new(None); - let (tx_a, mut rx_a) = mpsc::unbounded_channel::(); + let (tx_a, rx_a) = flume::unbounded::(); let handle_a = tokio::spawn(async move { shell_a .run( @@ -3528,7 +3526,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] let a_child_pid = time::timeout(Duration::from_secs(5), async { loop { let chunk = rx_a - .recv() + .recv_async() .await .expect("run A ended before printing its child pid"); a_output.push_str(&chunk); @@ -3852,7 +3850,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] async fn read_output_stops_when_cancelled_before_pipe_eof() { let (reader, _writer) = pipe_to_files("test").expect("test pipe should be created"); let cancel = CancellationToken::new(); - let (activity_tx, _activity_rx) = mpsc::channel(1); + let (activity_tx, _activity_rx) = flume::bounded(1); let handle = tokio::spawn(read_output(reader, None, cancel.clone(), activity_tx)); time::sleep(Duration::from_millis(10)).await; @@ -3867,8 +3865,8 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] #[cfg(unix)] #[tokio::test(flavor = "multi_thread")] async fn execute_shell_streams_separates_stdout_and_stderr() { - let (stdout_tx, mut stdout_rx) = mpsc::unbounded_channel::(); - let (stderr_tx, mut stderr_rx) = mpsc::unbounded_channel::(); + let (stdout_tx, stdout_rx) = flume::unbounded::(); + let (stderr_tx, stderr_rx) = flume::unbounded::(); let options = ShellExecuteOptions { command: "echo out; echo err 1>&2".to_string(), ..Default::default() @@ -3881,11 +3879,11 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] assert!(!result.cancelled); let mut stdout = Vec::new(); - while let Some(chunk) = stdout_rx.recv().await { + while let Ok(chunk) = stdout_rx.recv_async().await { stdout.extend_from_slice(&chunk); } let mut stderr = Vec::new(); - while let Some(chunk) = stderr_rx.recv().await { + while let Ok(chunk) = stderr_rx.recv_async().await { stderr.extend_from_slice(&chunk); } assert_eq!(stdout, b"out\n"); @@ -3914,7 +3912,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] #[cfg(unix)] #[tokio::test(flavor = "multi_thread")] async fn powershell_env_reference_survives_brush_expansion() { - let (tx, mut rx) = mpsc::unbounded_channel::(); + let (tx, rx) = flume::unbounded::(); let options = ShellExecuteOptions { command: "printf '%s' \"$env:SystemRoot\"".to_string(), ..Default::default() @@ -3926,7 +3924,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] assert_eq!(result.exit_code, Some(0)); let mut stdout = Vec::new(); - while let Some(chunk) = rx.recv().await { + while let Ok(chunk) = rx.recv_async().await { stdout.extend_from_slice(&chunk); } assert_eq!(stdout, b"$env:SystemRoot"); @@ -3938,7 +3936,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] #[cfg(unix)] #[tokio::test(flavor = "multi_thread")] async fn user_env_assignment_shadows_powershell_fallback() { - let (tx, mut rx) = mpsc::unbounded_channel::(); + let (tx, rx) = flume::unbounded::(); let options = ShellExecuteOptions { command: "env=prod; printf '%s' \"$env:8080\"".to_string(), ..Default::default() @@ -3950,7 +3948,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] assert_eq!(result.exit_code, Some(0)); let mut stdout = Vec::new(); - while let Some(chunk) = rx.recv().await { + while let Ok(chunk) = rx.recv_async().await { stdout.extend_from_slice(&chunk); } assert_eq!(stdout, b"prod:8080"); @@ -4026,7 +4024,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] #[cfg(unix)] #[tokio::test(flavor = "multi_thread")] async fn nohup_background_captures_operand_pid() { - let (tx, mut rx) = mpsc::unbounded_channel::(); + let (tx, rx) = flume::unbounded::(); let options = ShellExecuteOptions { command: "nohup /bin/sh -c 'exit 0' >/dev/null 2>&1 & pid=$!; printf 'pid=%s\n' \ \"$pid\"; test -n \"$pid\"" @@ -4041,7 +4039,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] assert!(!result.timed_out); let mut out = String::new(); - while let Some(chunk) = rx.recv().await { + while let Ok(chunk) = rx.recv_async().await { out.push_str(&chunk); } let pid = out @@ -4055,14 +4053,14 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] /// and exit code 125 (a nohup-level error, distinct from any command code). #[tokio::test(flavor = "multi_thread")] async fn nohup_builtin_without_command_reports_missing_operand() { - let (tx, mut rx) = mpsc::unbounded_channel::(); + let (tx, rx) = flume::unbounded::(); let options = ShellExecuteOptions { command: "nohup".to_string(), ..Default::default() }; let result = execute_shell(options, Some(tx), CancelToken::default()) .await .expect("execute should succeed"); assert_eq!(result.exit_code, Some(125)); let mut out = String::new(); - while let Some(chunk) = rx.recv().await { + while let Ok(chunk) = rx.recv_async().await { out.push_str(&chunk); } assert!( @@ -4096,7 +4094,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] let probe = "import signal,sys; sys.stdout.write('IGN' if \ signal.getsignal(signal.SIGHUP)==signal.SIG_IGN else 'DFL')"; - let (tx, mut rx) = mpsc::unbounded_channel::(); + let (tx, rx) = flume::unbounded::(); let options = ShellExecuteOptions { command: format!("nohup python3 -c \"{probe}\""), ..Default::default() @@ -4106,7 +4104,7 @@ replace = [{ pattern = "^.+$", replacement = "PWD" }] .expect("execute should succeed"); assert_eq!(result.exit_code, Some(0)); let mut out = String::new(); - while let Some(chunk) = rx.recv().await { + while let Ok(chunk) = rx.recv_async().await { out.push_str(&chunk); } assert!( diff --git a/crates/pi-uu-grep/Cargo.toml b/crates/pi-uu-grep/Cargo.toml index 2e3f3aa48..e27b2f73d 100644 --- a/crates/pi-uu-grep/Cargo.toml +++ b/crates/pi-uu-grep/Cargo.toml @@ -20,3 +20,6 @@ grep-regex = "0.1" grep-searcher = "0.1" ignore = "0.4" globset = "0.4" + +[dev-dependencies] +parking_lot.workspace = true diff --git a/crates/pi-uu-grep/src/lib.rs b/crates/pi-uu-grep/src/lib.rs index e2eb756f0..40ca9ecf6 100644 --- a/crates/pi-uu-grep/src/lib.rs +++ b/crates/pi-uu-grep/src/lib.rs @@ -768,9 +768,10 @@ mod tests { use std::{ collections::HashMap, io::Cursor, - sync::{Arc, Mutex, atomic::AtomicBool}, + sync::{Arc, atomic::AtomicBool}, }; + use parking_lot::Mutex; use pi_uutils_ctx::{ScopeIo, scope}; use super::*; @@ -780,7 +781,7 @@ mod tests { impl Write for SharedBuf { fn write(&mut self, buf: &[u8]) -> io::Result { - self.0.lock().expect("buffer lock").extend_from_slice(buf); + self.0.lock().extend_from_slice(buf); Ok(buf.len()) } @@ -809,10 +810,8 @@ mod tests { .map(OsString::from) .collect(); let code = scope(io, || run(argv)); - let stdout = - String::from_utf8(out.lock().expect("stdout lock").clone()).expect("utf8 stdout"); - let stderr = - String::from_utf8(err.lock().expect("stderr lock").clone()).expect("utf8 stderr"); + let stdout = String::from_utf8(out.lock().clone()).expect("utf8 stdout"); + let stderr = String::from_utf8(err.lock().clone()).expect("utf8 stderr"); (code, stdout, stderr) } diff --git a/crates/pi-walker/Cargo.toml b/crates/pi-walker/Cargo.toml index 17e78d88f..858680fa1 100644 --- a/crates/pi-walker/Cargo.toml +++ b/crates/pi-walker/Cargo.toml @@ -13,6 +13,7 @@ workspace = true dashmap.workspace = true globset.workspace = true ignore.workspace = true +parking_lot.workspace = true rayon.workspace = true [target.'cfg(unix)'.dependencies] diff --git a/crates/pi-walker/src/cache.rs b/crates/pi-walker/src/cache.rs index 37b93f53b..82afdea6e 100644 --- a/crates/pi-walker/src/cache.rs +++ b/crates/pi-walker/src/cache.rs @@ -4,12 +4,13 @@ use std::{ borrow::Cow, fmt, path::{Path, PathBuf}, - sync::{Arc, LazyLock, Mutex}, + sync::{Arc, LazyLock}, time::{Duration, Instant}, }; use dashmap::DashMap; use ignore::{ParallelVisitor, ParallelVisitorBuilder, WalkBuilder, WalkState}; +use parking_lot::Mutex; use rayon::{ThreadPool, prelude::*}; use crate::{ @@ -280,7 +281,6 @@ fn build_walker_for_options_inner( if let Some(pruned_dirs) = &pruned_dirs && pruned_dirs .lock() - .expect("pruned directory lock poisoned") .iter() .any(|dir| entry.path().starts_with(dir)) { @@ -384,11 +384,7 @@ impl Drop for EntryVisitor<'_, H> { return; } let entries = std::mem::take(&mut self.entries); - self - .shared_entries - .lock() - .expect("entry collection lock poisoned") - .push(entries); + self.shared_entries.lock().push(entries); } } @@ -401,7 +397,7 @@ where if self.visited == 0 || self.visited >= 128 { self.visited = 0; if let Err(err) = (self.heartbeat)() { - *self.error.lock().expect("error lock poisoned") = Some(err.to_string()); + *self.error.lock() = Some(err.to_string()); return WalkState::Quit; } } @@ -489,18 +485,12 @@ where heartbeat().map_err(|err| WalkError::Interrupted(err.to_string()))?; builder.build_parallel().visit(&mut visitor_builder); - let walk_error = error.lock().expect("error lock poisoned").take(); + let walk_error = error.lock().take(); if let Some(error) = walk_error { return Err(WalkError::Interrupted(error)); } - entries.extend( - shared_entries - .lock() - .expect("entry collection lock poisoned") - .drain(..) - .flatten(), - ); + entries.extend(shared_entries.lock().drain(..).flatten()); entries.sort_unstable_by(|a, b| a.path.cmp(&b.path)); Ok(EntryScan::Entries(CollectedEntries { entries, cache_age_ms: 0 })) } diff --git a/crates/pi-walker/src/lib.rs b/crates/pi-walker/src/lib.rs index fa9b910af..e2100f02b 100644 --- a/crates/pi-walker/src/lib.rs +++ b/crates/pi-walker/src/lib.rs @@ -17,7 +17,7 @@ use std::{ hash::{Hash, Hasher}, io, path::{Path, PathBuf}, - sync::{Arc, Mutex}, + sync::Arc, }; pub use cache::{ @@ -26,6 +26,7 @@ pub use cache::{ parallel_for_each, resolve_search_path, should_parallelize, should_skip_path, walk_workers, }; use globset::{GlobBuilder, GlobSet, GlobSetBuilder}; +use parking_lot::Mutex; const HEARTBEAT_INTERVAL: usize = 128; @@ -2307,10 +2308,7 @@ where WalkControl::Quit => return Ok(WalkStatus::Stopped), WalkControl::SkipDescend => { if collected.file_type == FileType::Dir { - pruned_dirs - .lock() - .expect("pruned directory lock poisoned") - .push(entry.path().to_path_buf()); + pruned_dirs.lock().push(entry.path().to_path_buf()); } }, WalkControl::Continue => {}, @@ -2361,11 +2359,7 @@ fn ignore_error_to_io(error: &ignore::Error) -> io::Error { } fn is_pruned_path(path: &Path, pruned_dirs: &Arc>>) -> bool { - pruned_dirs - .lock() - .expect("pruned directory lock poisoned") - .iter() - .any(|dir| path.starts_with(dir)) + pruned_dirs.lock().iter().any(|dir| path.starts_with(dir)) } /// Return whether [`WalkDetail::Full`] provides file sizes without per-entry diff --git a/crates/vendor/uu-sort/Cargo.toml b/crates/vendor/uu-sort/Cargo.toml index dce9d3891..a50e57f39 100644 --- a/crates/vendor/uu-sort/Cargo.toml +++ b/crates/vendor/uu-sort/Cargo.toml @@ -36,6 +36,7 @@ uucore = { version = "0.8.0", features = [ "i18n-collator", "i18n-datetime", ] } +flume = { workspace = true } pi-uutils-ctx = { path = "../../pi-uutils-ctx" } [target.'cfg(unix)'.dependencies] diff --git a/crates/vendor/uu-sort/src/check.rs b/crates/vendor/uu-sort/src/check.rs index e30388970..4e3978ca1 100644 --- a/crates/vendor/uu-sort/src/check.rs +++ b/crates/vendor/uu-sort/src/check.rs @@ -5,15 +5,9 @@ //! Check if a file is ordered -use std::{ - cmp::Ordering, - ffi::OsStr, - io::Read, - iter, - sync::mpsc::{Receiver, SyncSender, sync_channel}, - thread, -}; +use std::{cmp::Ordering, ffi::OsStr, io::Read, iter, thread}; +use flume::{Receiver, Sender}; use itertools::Itertools; use uucore::error::UResult; @@ -39,8 +33,8 @@ pub fn check(path: &OsStr, settings: &GlobalSettings) -> UResult<()> { Ordering::Equal }; let file = open(path)?; - let (recycled_sender, recycled_receiver) = sync_channel(2); - let (loaded_sender, loaded_receiver) = sync_channel(2); + let (recycled_sender, recycled_receiver) = flume::bounded(2); + let (loaded_sender, loaded_receiver) = flume::bounded(2); thread::spawn({ let settings = settings.clone(); move || reader(file, &recycled_receiver, &loaded_sender, &settings) @@ -57,7 +51,7 @@ pub fn check(path: &OsStr, settings: &GlobalSettings) -> UResult<()> { let mut prev_chunk: Option = None; let mut line_idx = 0; - for chunk in loaded_receiver { + while let Ok(chunk) = loaded_receiver.recv() { line_idx += 1; if let Some(prev_chunk) = prev_chunk.take() { // Check if the first element of the new chunk is greater than the last @@ -105,11 +99,11 @@ pub fn check(path: &OsStr, settings: &GlobalSettings) -> UResult<()> { fn reader( mut file: Box, receiver: &Receiver, - sender: &SyncSender, + sender: &Sender, settings: &GlobalSettings, ) -> UResult<()> { let mut carry_over = vec![]; - for recycled_chunk in receiver { + while let Ok(recycled_chunk) = receiver.recv() { let should_continue = chunks::read( sender, recycled_chunk, diff --git a/crates/vendor/uu-sort/src/chunks.rs b/crates/vendor/uu-sort/src/chunks.rs index 3acbd71aa..a273f1f22 100644 --- a/crates/vendor/uu-sort/src/chunks.rs +++ b/crates/vendor/uu-sort/src/chunks.rs @@ -12,9 +12,9 @@ use std::{ io::{ErrorKind, Read}, ops::Range, - sync::mpsc::SyncSender, }; +use flume::Sender; use memchr::memchr_iter; use self_cell::self_cell; use uucore::error::{UResult, USimpleError}; @@ -180,7 +180,7 @@ impl RecycledChunk { /// * `settings`: The global settings. #[allow(clippy::too_many_arguments)] pub fn read( - sender: &SyncSender, + sender: &Sender, recycled_chunk: RecycledChunk, max_buffer_size: Option, carry_over: &mut Vec, diff --git a/crates/vendor/uu-sort/src/ext_sort/threaded.rs b/crates/vendor/uu-sort/src/ext_sort/threaded.rs index 211309ec5..54d240b56 100644 --- a/crates/vendor/uu-sort/src/ext_sort/threaded.rs +++ b/crates/vendor/uu-sort/src/ext_sort/threaded.rs @@ -11,10 +11,10 @@ use std::{ fs::File, io::{Read, Write}, path::PathBuf, - sync::mpsc::{Receiver, SyncSender}, thread, }; +use flume::{Receiver, Sender}; use itertools::Itertools; use uucore::error::{UResult, strip_errno}; @@ -44,8 +44,8 @@ pub fn ext_sort( output: Output, tmp_dir: &mut TmpDirWrapper, ) -> UResult<()> { - let (sorted_sender, sorted_receiver) = std::sync::mpsc::sync_channel(1); - let (recycled_sender, recycled_receiver) = std::sync::mpsc::sync_channel(1); + let (sorted_sender, sorted_receiver) = flume::bounded(1); + let (recycled_sender, recycled_receiver) = flume::bounded(1); thread::spawn({ let settings = settings.clone(); move || sorter(&recycled_receiver, &sorted_sender, &settings) @@ -105,7 +105,7 @@ fn reader_writer< files: F, settings: &GlobalSettings, receiver: &Receiver, - sender: SyncSender, + sender: Sender, output: Output, tmp_dir: &mut TmpDirWrapper, ) -> UResult<()> { @@ -177,7 +177,7 @@ fn reader_writer< } /// The function that is executed on the sorter thread. -fn sorter(receiver: &Receiver, sender: &SyncSender, settings: &GlobalSettings) { +fn sorter(receiver: &Receiver, sender: &Sender, settings: &GlobalSettings) { while let Ok(mut payload) = receiver.recv() { payload.with_dependent_mut(|_, contents| { sort_by(&mut contents.lines, settings, &contents.line_data); @@ -210,7 +210,7 @@ fn read_write_loop( buffer_size: usize, settings: &GlobalSettings, receiver: &Receiver, - sender: SyncSender, + sender: Sender, ) -> UResult> { let mut file = files.next().unwrap()?; diff --git a/crates/vendor/uu-sort/src/merge.rs b/crates/vendor/uu-sort/src/merge.rs index 3a41d883e..30b052160 100644 --- a/crates/vendor/uu-sort/src/merge.rs +++ b/crates/vendor/uu-sort/src/merge.rs @@ -23,11 +23,11 @@ use std::{ path::PathBuf, process::{Child, ChildStdin, ChildStdout, Command, Stdio}, rc::Rc, - sync::mpsc::{Receiver, Sender, SyncSender, channel, sync_channel}, thread::{self, JoinHandle}, }; use compare::Compare; +use flume::{Receiver, Sender}; use uucore::error::{FromIo, UResult}; use crate::{ @@ -177,11 +177,11 @@ fn merge_without_limit>>( files: F, settings: &GlobalSettings, ) -> UResult> { - let (request_sender, request_receiver) = channel(); + let (request_sender, request_receiver) = flume::unbounded(); let mut reader_files = Vec::with_capacity(files.size_hint().0); let mut loaded_receivers = Vec::with_capacity(files.size_hint().0); for (file_number, file) in files.enumerate() { - let (sender, receiver) = sync_channel(2); + let (sender, receiver) = flume::bounded(2); loaded_receivers.push(receiver); reader_files.push(Some(ReaderFile { file: file?, sender, carry_over: vec![] })); // Send the initial chunk to trigger a read for each file @@ -227,7 +227,7 @@ fn merge_without_limit>>( /// The struct on the reader thread representing an input file struct ReaderFile { file: M, - sender: SyncSender, + sender: Sender, carry_over: Vec, } @@ -238,7 +238,7 @@ fn reader( settings: &GlobalSettings, separator: u8, ) -> UResult<()> { - for (file_idx, recycled_chunk) in recycled_receiver { + while let Ok((file_idx, recycled_chunk)) = recycled_receiver.recv() { if let Some(ReaderFile { file, sender, carry_over }) = &mut files[file_idx] { let should_continue = chunks::read( sender,