Skip to main content

skimmer/helper/
item_reader.rs

1/// helper for turn a BufRead into a skim stream
2use std::env;
3use std::error::Error;
4use std::io::{BufRead, BufReader};
5use std::process::{Child, Command, Stdio};
6use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
7use std::sync::Arc;
8use std::thread;
9
10use crossbeam::channel::{bounded, Receiver, Sender};
11use regex::Regex;
12
13use crate::field::FieldRange;
14use crate::helper::item::DefaultSkimItem;
15use crate::reader::CommandCollector;
16use crate::{SkimItem, SkimItemReceiver, SkimItemSender};
17
18const CMD_CHANNEL_SIZE: usize = 1024;
19const ITEM_CHANNEL_SIZE: usize = 10240;
20const DELIMITER_STR: &str = r"[\t\n ]+";
21const READ_BUFFER_SIZE: usize = 1024;
22
23pub enum CollectorInput {
24    Pipe(Box<dyn BufRead + Send>),
25    Command(String),
26}
27
28#[derive(Debug)]
29pub struct SkimItemReaderOption {
30    buf_size: usize,
31    use_ansi_color: bool,
32    transform_fields: Vec<FieldRange>,
33    matching_fields: Vec<FieldRange>,
34    delimiter: Regex,
35    line_ending: u8,
36    show_error: bool,
37}
38
39impl Default for SkimItemReaderOption {
40    fn default() -> Self {
41        Self {
42            buf_size: READ_BUFFER_SIZE,
43            line_ending: b'\n',
44            use_ansi_color: false,
45            transform_fields: Vec::new(),
46            matching_fields: Vec::new(),
47            delimiter: Regex::new(DELIMITER_STR).unwrap(),
48            show_error: false,
49        }
50    }
51}
52
53impl SkimItemReaderOption {
54    pub fn buf_size(mut self, buf_size: usize) -> Self {
55        self.buf_size = buf_size;
56        self
57    }
58
59    pub fn line_ending(mut self, line_ending: u8) -> Self {
60        self.line_ending = line_ending;
61        self
62    }
63
64    pub fn ansi(mut self, enable: bool) -> Self {
65        self.use_ansi_color = enable;
66        self
67    }
68
69    pub fn delimiter(mut self, delimiter: &str) -> Self {
70        if !delimiter.is_empty() {
71            self.delimiter = Regex::new(delimiter).unwrap_or_else(|_| Regex::new(DELIMITER_STR).unwrap());
72        }
73        self
74    }
75
76    pub fn with_nth(mut self, with_nth: &[String]) -> Self {
77        self.transform_fields = with_nth.iter().filter_map(|s| FieldRange::from_str(s)).collect();
78        self
79    }
80
81    pub fn transform_fields(mut self, transform_fields: Vec<FieldRange>) -> Self {
82        self.transform_fields = transform_fields;
83        self
84    }
85
86    pub fn nth(mut self, nth: &[String]) -> Self {
87        self.matching_fields = nth.iter().filter_map(|s| FieldRange::from_str(s)).collect();
88        self
89    }
90
91    pub fn matching_fields(mut self, matching_fields: Vec<FieldRange>) -> Self {
92        self.matching_fields = matching_fields;
93        self
94    }
95
96    pub fn read0(mut self, enable: bool) -> Self {
97        if enable {
98            self.line_ending = b'\0';
99        } else {
100            self.line_ending = b'\n';
101        }
102        self
103    }
104
105    pub fn show_error(mut self, show_error: bool) -> Self {
106        self.show_error = show_error;
107        self
108    }
109
110    pub fn build(self) -> Self {
111        self
112    }
113
114    pub fn is_simple(&self) -> bool {
115        debug!("transorm {:?}", self.transform_fields);
116        !self.use_ansi_color && self.matching_fields.is_empty() && self.transform_fields.is_empty()
117    }
118}
119
120pub struct SkimItemReader {
121    option: Arc<SkimItemReaderOption>,
122}
123
124impl Default for SkimItemReader {
125    fn default() -> Self {
126        Self {
127            option: Arc::new(Default::default()),
128        }
129    }
130}
131
132impl SkimItemReader {
133    pub fn new(option: SkimItemReaderOption) -> Self {
134        Self {
135            option: Arc::new(option),
136        }
137    }
138
139    pub fn option(mut self, option: SkimItemReaderOption) -> Self {
140        self.option = Arc::new(option);
141        self
142    }
143}
144
145impl SkimItemReader {
146    pub fn of_bufread(&self, source: impl BufRead + Send + 'static) -> SkimItemReceiver {
147        if self.option.is_simple() {
148            debug!("Is simple");
149            self.raw_bufread(source)
150        } else {
151            self.read_and_collect_from_command(Arc::new(AtomicUsize::new(0)), CollectorInput::Pipe(Box::new(source)))
152                .0
153        }
154    }
155
156    /// helper: convert bufread into SkimItemReceiver
157    fn raw_bufread(&self, mut source: impl BufRead + Send + 'static) -> SkimItemReceiver {
158        let (tx_item, rx_item): (SkimItemSender, SkimItemReceiver) = bounded(self.option.buf_size);
159        let line_ending = self.option.line_ending;
160        thread::spawn(move || {
161            let mut buffer = Vec::with_capacity(1024);
162            loop {
163                buffer.clear();
164                // start reading
165                match source.read_until(line_ending, &mut buffer) {
166                    Ok(n) => {
167                        if n == 0 {
168                            break;
169                        }
170
171                        if buffer.ends_with(b"\r\n") {
172                            buffer.pop();
173                            buffer.pop();
174                        } else if buffer.ends_with(b"\n") || buffer.ends_with(b"\0") {
175                            buffer.pop();
176                        }
177
178                        let string = String::from_utf8_lossy(&buffer);
179                        let result = tx_item.send(Arc::new(string.into_owned()));
180                        if result.is_err() {
181                            break;
182                        }
183                    }
184                    Err(_err) => {} // String not UTF8 or other error, skip.
185                }
186            }
187        });
188        rx_item
189    }
190
191    /// components_to_stop == 0 => all the threads have been stopped
192    /// return (channel_for_receive_item, channel_to_stop_command)
193    fn read_and_collect_from_command(
194        &self,
195        components_to_stop: Arc<AtomicUsize>,
196        input: CollectorInput,
197    ) -> (Receiver<Arc<dyn SkimItem>>, Sender<i32>) {
198        let (command, mut source) = match input {
199            CollectorInput::Pipe(pipe) => (None, pipe),
200            CollectorInput::Command(cmd) => get_command_output(&cmd).expect("command not found"),
201        };
202
203        let (tx_interrupt, rx_interrupt) = bounded(CMD_CHANNEL_SIZE);
204        let (tx_item, rx_item): (SkimItemSender, SkimItemReceiver) = bounded(ITEM_CHANNEL_SIZE);
205
206        let started = Arc::new(AtomicBool::new(false));
207        let started_clone = started.clone();
208        let components_to_stop_clone = components_to_stop.clone();
209        let tx_item_clone = tx_item.clone();
210        let send_error = self.option.show_error;
211        // listening to close signal and kill command if needed
212        thread::spawn(move || {
213            debug!("collector: command killer start");
214            components_to_stop_clone.fetch_add(1, Ordering::SeqCst);
215            started_clone.store(true, Ordering::SeqCst); // notify parent that it is started
216
217            let _ = rx_interrupt.recv(); // block waiting
218            if let Some(mut child) = command {
219                // clean up resources
220                let _ = child.kill();
221                let _ = child.wait();
222
223                if send_error {
224                    let has_error = child
225                        .try_wait()
226                        .map(|os| os.map(|s| !s.success()).unwrap_or(true))
227                        .unwrap_or(false);
228                    if has_error {
229                        let output = child.wait_with_output().expect("could not retrieve error message");
230                        for line in String::from_utf8_lossy(&output.stderr).lines() {
231                            let _ = tx_item_clone.send(Arc::new(line.to_string()));
232                        }
233                    }
234                }
235            }
236
237            components_to_stop_clone.fetch_sub(1, Ordering::SeqCst);
238            debug!("collector: command killer stop");
239        });
240
241        while !started.load(Ordering::SeqCst) {
242            // busy waiting for the thread to start. (components_to_stop is added)
243        }
244
245        let started = Arc::new(AtomicBool::new(false));
246        let started_clone = started.clone();
247        let tx_interrupt_clone = tx_interrupt.clone();
248        let option = self.option.clone();
249        thread::spawn(move || {
250            debug!("collector: command collector start");
251            components_to_stop.fetch_add(1, Ordering::SeqCst);
252            started_clone.store(true, Ordering::SeqCst); // notify parent that it is started
253
254            let mut buffer = Vec::with_capacity(option.buf_size);
255            loop {
256                buffer.clear();
257
258                // start reading
259                match source.read_until(option.line_ending, &mut buffer) {
260                    Ok(n) => {
261                        if n == 0 {
262                            break;
263                        }
264
265                        if buffer.ends_with(b"\r\n") {
266                            buffer.pop();
267                            buffer.pop();
268                        } else if buffer.ends_with(b"\n") || buffer.ends_with(b"\0") {
269                            buffer.pop();
270                        }
271
272                        let line = String::from_utf8_lossy(&buffer).to_string();
273
274                        let raw_item = DefaultSkimItem::new(
275                            line,
276                            option.use_ansi_color,
277                            &option.transform_fields,
278                            &option.matching_fields,
279                            &option.delimiter,
280                        );
281
282                        match tx_item.send(Arc::new(raw_item)) {
283                            Ok(_) => {}
284                            Err(_) => {
285                                debug!("collector: failed to send item, quit");
286                                break;
287                            }
288                        }
289                    }
290                    Err(_err) => {} // String not UTF8 or other error, skip.
291                }
292            }
293
294            let _ = tx_interrupt_clone.send(1); // ensure the waiting thread will exit
295            components_to_stop.fetch_sub(1, Ordering::SeqCst);
296            debug!("collector: command collector stop");
297        });
298
299        while !started.load(Ordering::SeqCst) {
300            // busy waiting for the thread to start. (components_to_stop is added)
301        }
302
303        (rx_item, tx_interrupt)
304    }
305}
306
307impl CommandCollector for SkimItemReader {
308    fn invoke(&mut self, cmd: &str, components_to_stop: Arc<AtomicUsize>) -> (SkimItemReceiver, Sender<i32>) {
309        self.read_and_collect_from_command(components_to_stop, CollectorInput::Command(cmd.to_string()))
310    }
311}
312
313type CommandOutput = (Option<Child>, Box<dyn BufRead + Send>);
314
315fn get_command_output(cmd: &str) -> Result<CommandOutput, Box<dyn Error>> {
316    let shell = env::var("SHELL").unwrap_or_else(|_| "sh".to_string());
317    let mut command: Child = Command::new(shell)
318        .arg("-c")
319        .arg(cmd)
320        .stdout(Stdio::piped())
321        .stderr(Stdio::piped())
322        .spawn()?;
323
324    let stdout = command
325        .stdout
326        .take()
327        .ok_or_else(|| "command output: unwrap failed".to_owned())?;
328
329    Ok((Some(command), Box::new(BufReader::new(stdout))))
330}