diff --git a/bench.sh b/bench.sh index 48424b3c..868ec7ff 100755 --- a/bench.sh +++ b/bench.sh @@ -182,6 +182,8 @@ for RUN in $(seq 1 $RUNS); do # Unset HISTFILE in the tmux session to prevent command from appearing in shell history tmux send-keys -t "$SESSION_NAME" "unset HISTFILE" Enter + tmux send-keys -t "$SESSION_NAME" "unset FZF_DEFAULT_OPTS" Enter + tmux send-keys -t "$SESSION_NAME" "unset SKIM_DEFAULT_OPTIONS" Enter sleep 0.1 # Prepare to capture the start time as close to data ingestion as possible diff --git a/examples/cmd_collector.rs b/examples/cmd_collector.rs index d820439c..6d861e8d 100644 --- a/examples/cmd_collector.rs +++ b/examples/cmd_collector.rs @@ -20,9 +20,13 @@ impl CommandCollector for BasicCmdCollector { fn invoke(&mut self, _cmd: &str, _components_to_stop: Arc) -> (SkimItemReceiver, Sender) { let (tx, rx) = unbounded(); let (tx_interrupt, _rx_interrupt) = unbounded(); + let mut batch = Vec::new(); while let Some(value) = self.items.pop() { let item = BasicSkimItem { value }; - tx.send(Arc::from(item) as Arc).unwrap(); + batch.push(Arc::from(item) as Arc); + } + if !batch.is_empty() { + tx.send(batch).unwrap(); } (rx, tx_interrupt) diff --git a/examples/custom_item.rs b/examples/custom_item.rs index b77fcdcf..6b677cb6 100644 --- a/examples/custom_item.rs +++ b/examples/custom_item.rs @@ -28,15 +28,17 @@ fn main() { .unwrap(); let (tx_item, rx_item): (SkimItemSender, SkimItemReceiver) = unbounded(); - let _ = tx_item.send(Arc::new(MyItem { - inner: "color aaaa".to_string(), - })); - let _ = tx_item.send(Arc::new(MyItem { - inner: "bbbb".to_string(), - })); - let _ = tx_item.send(Arc::new(MyItem { - inner: "ccc".to_string(), - })); + let _ = tx_item.send(vec![ + Arc::new(MyItem { + inner: "color aaaa".to_string(), + }) as Arc, + Arc::new(MyItem { + inner: "bbbb".to_string(), + }) as Arc, + Arc::new(MyItem { + inner: "ccc".to_string(), + }) as Arc, + ]); drop(tx_item); // so that skim could know when to stop waiting for more items. let selected_items = Skim::run_with(options, Some(rx_item)) diff --git a/examples/downcast.rs b/examples/downcast.rs index d75f2242..1931ed2d 100644 --- a/examples/downcast.rs +++ b/examples/downcast.rs @@ -38,20 +38,20 @@ pub fn main() { let (tx, rx): (SkimItemSender, SkimItemReceiver) = unbounded(); - tx.send(Arc::new(Item { - text: "a".to_string(), - index: 0, - })) - .unwrap(); - tx.send(Arc::new(Item { - text: "b".to_string(), - index: 1, - })) - .unwrap(); - tx.send(Arc::new(Item { - text: "c".to_string(), - index: 2, - })) + tx.send(vec![ + Arc::new(Item { + text: "a".to_string(), + index: 0, + }) as Arc, + Arc::new(Item { + text: "b".to_string(), + index: 1, + }) as Arc, + Arc::new(Item { + text: "c".to_string(), + index: 2, + }) as Arc, + ]) .unwrap(); drop(tx); diff --git a/examples/receiver_multi.rs b/examples/receiver_multi.rs index c10b1329..f1d6c4b4 100644 --- a/examples/receiver_multi.rs +++ b/examples/receiver_multi.rs @@ -3,10 +3,12 @@ use std::sync::Arc; use skim::prelude::*; fn main() { - let (sender, receiver) = unbounded::>(); + let (sender, receiver): (SkimItemSender, SkimItemReceiver) = unbounded(); + let mut batch = Vec::new(); for num in 1..=8 { - sender.send(Arc::new(format!("Option {num}"))).unwrap(); + batch.push(Arc::new(format!("Option {num}")) as Arc); } + sender.send(batch).unwrap(); drop(sender); // bug replicates even without this let _ = Skim::run_with( diff --git a/src/bin/main.rs b/src/bin/main.rs index ce3a3179..9e432a51 100644 --- a/src/bin/main.rs +++ b/src/bin/main.rs @@ -331,8 +331,8 @@ pub fn filter(bin_option: &BinOptions, options: &SkimOptions, source: Option = items diff --git a/src/helper/item_reader.rs b/src/helper/item_reader.rs index d5897dc8..4c99badf 100644 --- a/src/helper/item_reader.rs +++ b/src/helper/item_reader.rs @@ -6,9 +6,10 @@ use std::process::{Child, Command, Stdio}; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::thread; +use std::time::{Duration, Instant}; use regex::Regex; -use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender, unbounded_channel}; +use tokio::sync::mpsc::{UnboundedSender, unbounded_channel}; use crate::field::FieldRange; use crate::helper::item::DefaultSkimItem; @@ -17,6 +18,8 @@ use crate::{SkimItem, SkimItemReceiver, SkimItemSender, SkimOptions}; const DELIMITER_STR: &str = r"[\t\n ]+"; const READ_BUFFER_SIZE: usize = 1024; +const ITEMS_BUFFER_SIZE: usize = 128; +const SEND_TIMEOUT_MS: u64 = 100; // Send items if we haven't sent anything in 100ms pub enum CollectorInput { Pipe(Box), @@ -191,50 +194,90 @@ impl SkimItemReader { } } - /// helper: convert bufread into SkimItemReceiver - fn raw_bufread(&self, mut source: impl BufRead + Send + 'static) -> SkimItemReceiver { - let (tx_item, rx_item): (SkimItemSender, SkimItemReceiver) = unbounded_channel(); - let line_ending = self.option.line_ending; - let use_ansi = self.option.use_ansi_color; - let delimiter = self.option.delimiter.clone(); - thread::spawn(move || { - let mut buffer = Vec::with_capacity(1024); - let mut idx = 0; - loop { - buffer.clear(); - // start reading - match source.read_until(line_ending, &mut buffer) { - Ok(n) => { - if n == 0 { - break; - } + /// Helper function that contains the common logic for reading lines from a BufRead source + /// and converting them into SkimItems. + fn read_lines_into_items( + mut source: impl BufRead + Send + 'static, + tx_item: SkimItemSender, + option: Arc, + transform_fields: Vec, + matching_fields: Vec, + ) { + let mut buffer = Vec::with_capacity(option.buf_size); + let mut line_idx = 0; + let mut items_to_send = Vec::with_capacity(ITEMS_BUFFER_SIZE); + let mut last_send_time = Instant::now(); + let send_timeout = Duration::from_millis(SEND_TIMEOUT_MS); - if buffer.ends_with(b"\r\n") { - buffer.pop(); - buffer.pop(); - } else if buffer.ends_with(b"\n") || buffer.ends_with(b"\0") { - buffer.pop(); - } + loop { + buffer.clear(); - let string = String::from_utf8_lossy(&buffer); - //let result = tx_item.send(Arc::new(string.into_owned())); - let result = tx_item.send(Arc::new(DefaultSkimItem::new( - string.to_string(), - use_ansi, - &[], - &[], - &delimiter, - idx, - ))); - if result.is_err() { - break; - } - idx += 1; + // start reading + match source.read_until(option.line_ending, &mut buffer) { + Ok(0) => break, + Ok(_) => { + // Strip line endings + if buffer.ends_with(b"\r\n") { + buffer.pop(); + buffer.pop(); + } else if buffer.ends_with(&[option.line_ending]) { + buffer.pop(); + } + + let line = String::from_utf8_lossy(&buffer).to_string(); + + trace!("got item {} with index {}", line.clone(), line_idx); + + let raw_item = DefaultSkimItem::new( + line, + option.use_ansi_color, + &transform_fields, + &matching_fields, + &option.delimiter, + line_idx, + ); + items_to_send.push(Arc::new(raw_item) as Arc); + + line_idx += 1; + } + Err(err) => { + trace!("Got {err:?} when reading, skipping"); + } // String not UTF8 or other error, skip. + } + + // Send batched items if buffer is full OR timeout has elapsed + let should_send = items_to_send.len() == ITEMS_BUFFER_SIZE + || (!items_to_send.is_empty() && last_send_time.elapsed() >= send_timeout); + + if should_send { + let batch = std::mem::replace(&mut items_to_send, Vec::with_capacity(ITEMS_BUFFER_SIZE)); + match tx_item.send(batch) { + Ok(_) => { + last_send_time = Instant::now(); + } + Err(e) => { + warn!("Failed to send items: {e:?}"); + break; } - Err(_err) => {} // String not UTF8 or other error, skip. } } + } + + // Send remaining items + if !items_to_send.is_empty() { + let _ = tx_item.send(items_to_send); + } + } + + /// helper: convert bufread into SkimItemReceiver + fn raw_bufread(&self, source: impl BufRead + Send + 'static) -> SkimItemReceiver { + let (tx_item, rx_item): (SkimItemSender, SkimItemReceiver) = unbounded_channel(); + let option = self.option.clone(); + + thread::spawn(move || { + Self::read_lines_into_items(source, tx_item, option, vec![], vec![]); }); + rx_item } @@ -244,15 +287,15 @@ impl SkimItemReader { &self, components_to_stop: Arc, input: CollectorInput, - ) -> (UnboundedReceiver>, UnboundedSender) { + ) -> (SkimItemReceiver, UnboundedSender) { let send_error = self.option.show_error; - let (command, mut source) = match input { + let (command, source) = match input { CollectorInput::Pipe(pipe) => (None, pipe), CollectorInput::Command(cmd) => get_command_output(&cmd, send_error).expect("command not found"), }; let (tx_interrupt, mut rx_interrupt) = unbounded_channel(); - let (tx_item, rx_item): (SkimItemSender, SkimItemReceiver) = unbounded_channel::>(); + let (tx_item, rx_item): (SkimItemSender, SkimItemReceiver) = unbounded_channel(); let started = Arc::new(AtomicBool::new(false)); let started_clone = started.clone(); @@ -278,9 +321,21 @@ impl SkimItemReader { if has_error { trace!("collector: sending error"); let output = child.wait_with_output().expect("could not retrieve error message"); - for line in String::from_utf8_lossy(&output.stderr).lines() { - let _ = tx_item_clone.send(Arc::new(line.to_string())); - } + let error_text = String::from_utf8_lossy(&output.stderr).to_string(); + let error_items: Vec> = error_text + .lines() + .map(|line| { + Arc::new(DefaultSkimItem::new( + line.to_string(), + false, + &[], + &[], + &Regex::new(DELIMITER_STR).unwrap(), + 0, + )) as Arc + }) + .collect(); + let _ = tx_item_clone.send(error_items); } } } @@ -297,54 +352,15 @@ impl SkimItemReader { let started_clone = started.clone(); let tx_interrupt_clone = tx_interrupt.clone(); let option = self.option.clone(); + let transform_fields = option.transform_fields.clone(); + let matching_fields = option.matching_fields.clone(); + thread::spawn(move || { debug!("collector: command collector start"); components_to_stop.fetch_add(1, Ordering::SeqCst); started_clone.store(true, Ordering::SeqCst); // notify parent that it is started - let mut buffer = Vec::with_capacity(option.buf_size); - let mut line_idx = 0; - loop { - buffer.clear(); - - // start reading - match source.read_until(option.line_ending, &mut buffer) { - Ok(0) => break, - Ok(_) => { - if buffer.ends_with(b"\r\n") { - buffer.pop(); - buffer.pop(); - } else if buffer.ends_with(&[option.line_ending]) { - buffer.pop(); - } - - let line = String::from_utf8_lossy(&buffer).to_string(); - - trace!("got item {} with index {} from command", line.clone(), line_idx); - - let raw_item = DefaultSkimItem::new( - line, - option.use_ansi_color, - &option.transform_fields, - &option.matching_fields, - &option.delimiter, - line_idx, - ); - - match tx_item.send(Arc::new(raw_item)) { - Ok(_) => {} - Err(_) => { - debug!("collector: failed to send item, quit"); - break; - } - } - line_idx += 1; - } - Err(err) => { - trace!("Got {err:?} when reading from command collector, skipping"); - } // String not UTF8 or other error, skip. - } - } + Self::read_lines_into_items(source, tx_item, option, transform_fields, matching_fields); let _ = tx_interrupt_clone.send(1); // ensure the waiting thread will exit components_to_stop.fetch_sub(1, Ordering::SeqCst); diff --git a/src/lib.rs b/src/lib.rs index e94a2c17..649f7042 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -328,9 +328,9 @@ pub trait Selector { //------------------------------------------------------------------------------ /// Sender for streaming items to skim -pub type SkimItemSender = UnboundedSender>; +pub type SkimItemSender = UnboundedSender>>; /// Receiver for streaming items to skim -pub type SkimItemReceiver = UnboundedReceiver>; +pub type SkimItemReceiver = UnboundedReceiver>>; /// Main entry point for running skim pub struct Skim {} @@ -399,7 +399,6 @@ impl Skim { let item_pool = app.item_pool.clone(); tokio::spawn(async move { - const BATCH: usize = 4096; // Smaller batches for more responsive updates loop { if reader_interrupt_rx .try_recv() @@ -408,13 +407,16 @@ impl Skim { debug!("stopping reader receiver thread"); break; } - let mut buf = Vec::with_capacity(BATCH); trace!("getting items"); - if item_rx.recv_many(&mut buf, BATCH).await > 0 { - item_pool.append(buf); - trace!("Got new items, len {}", item_pool.len()); - } else { - reader_done_clone.store(true, std::sync::atomic::Ordering::Relaxed); + match item_rx.recv().await { + Some(batch) => { + item_pool.append(batch); + trace!("Got new items, len {}", item_pool.len()); + } + None => { + reader_done_clone.store(true, std::sync::atomic::Ordering::Relaxed); + break; + } } } }); @@ -516,7 +518,7 @@ impl Skim { app.reader_timer = Instant::now(); } else if ! reader_done.load(std::sync::atomic::Ordering::Relaxed) { reader_done.store(true, std::sync::atomic::Ordering::Relaxed); - app.restart_matcher(true); + app.restart_matcher(false); } app.handle_event(&mut tui, &evt)?; } diff --git a/src/matcher.rs b/src/matcher.rs index 8e59a98f..7a35bc19 100644 --- a/src/matcher.rs +++ b/src/matcher.rs @@ -28,11 +28,10 @@ pub struct MatcherControl { impl Default for MatcherControl { fn default() -> Self { Self { - // Default to stopped=true so initial state indicates "no matcher running" stopped: Arc::new(AtomicBool::new(true)), - processed: Arc::new(AtomicUsize::new(0)), - matched: Arc::new(AtomicUsize::new(0)), - items: Arc::new(SpinLock::new(Vec::new())), + processed: Default::default(), + matched: Default::default(), + items: Default::default(), } } } diff --git a/src/reader.rs b/src/reader.rs index 6396ea17..b03316b7 100644 --- a/src/reader.rs +++ b/src/reader.rs @@ -88,7 +88,7 @@ impl Reader { } /// Starts the reader and returns a control handle - pub fn run(&mut self, app_tx: UnboundedSender>, cmd: &str) -> ReaderControl { + pub fn run(&mut self, app_tx: UnboundedSender>>, cmd: &str) -> ReaderControl { let components_to_stop: Arc = Arc::new(AtomicUsize::new(0)); let items = Arc::new(SpinLock::new(Vec::new())); @@ -113,7 +113,7 @@ impl Reader { fn collect_item( components_to_stop: Arc, mut rx_item: SkimItemReceiver, - app_tx: UnboundedSender>, + app_tx: UnboundedSender>>, ) -> UnboundedSender { let (tx_interrupt, mut rx_interrupt) = unbounded_channel(); @@ -128,8 +128,8 @@ fn collect_item( select! { new_item = rx_item.recv() => { match new_item { - Some(item) => { - let _ = app_tx.send(item); + Some(items) => { + let _ = app_tx.send(items); } None => break, } diff --git a/src/tui/app.rs b/src/tui/app.rs index 0f6c74ad..bb33766c 100644 --- a/src/tui/app.rs +++ b/src/tui/app.rs @@ -1013,6 +1013,7 @@ impl<'a> App<'a> { ToggleInteractive => { self.options.interactive = !self.options.interactive; self.input.switch_mode(); + self.restart_matcher(true); } ToggleOut => { self.item_list.toggle(); @@ -1101,18 +1102,7 @@ impl<'a> App<'a> { } let matcher_stopped = self.matcher_control.stopped(); - - // If a matcher is still running, don't start a new one - just mark pending - // This prevents race conditions where a killed matcher's empty results - // overwrite a successful matcher's results - if !matcher_stopped { - if force || self.item_pool.num_not_taken() > 0 { - self.pending_matcher_restart = true; - } - return; - } - - if force || self.pending_matcher_restart || self.item_pool.num_not_taken() > 0 { + if force || self.pending_matcher_restart || (matcher_stopped && self.item_pool.num_not_taken() > 0) { // Reset debounce timer on any restart to prevent interference self.last_matcher_restart = std::time::Instant::now(); self.pending_matcher_restart = false;