Skip to main content

pagers_core/
crawl.rs

1use std::collections::{HashSet, VecDeque};
2use std::io::{self, BufRead};
3use std::num::{NonZeroU16, NonZeroUsize};
4use std::path::{Path, PathBuf};
5use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
6use std::sync::{Arc, Mutex};
7
8use dua_core::Order;
9
10use crate::Cancellation;
11use crate::mincore::PageMap;
12use crate::mode::DisplayMode;
13use crate::ops::{FileRange, Op, Stats};
14
15#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
16#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
17pub enum Threads {
18    #[default]
19    All,
20    Exact(NonZeroU16),
21}
22
23impl Threads {
24    pub fn get(self) -> usize {
25        match self {
26            Self::All => std::thread::available_parallelism().map_or(1, NonZeroUsize::get),
27            Self::Exact(threads) => usize::from(threads.get()),
28        }
29    }
30
31    fn effective(self) -> usize {
32        self.get()
33            .min(std::thread::available_parallelism().map_or(1, NonZeroUsize::get))
34    }
35}
36
37impl From<u16> for Threads {
38    fn from(threads: u16) -> Self {
39        NonZeroU16::new(threads).map_or(Self::All, Self::Exact)
40    }
41}
42
43impl std::fmt::Display for Threads {
44    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
45        match self {
46            Self::All => f.write_str("0"),
47            Self::Exact(threads) => write!(f, "{threads}"),
48        }
49    }
50}
51
52impl std::str::FromStr for Threads {
53    type Err = std::num::ParseIntError;
54
55    fn from_str(value: &str) -> Result<Self, Self::Err> {
56        value.parse::<u16>().map(Self::from)
57    }
58}
59
60#[derive(Debug, Clone, PartialEq)]
61#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
62pub struct CrawlConfig {
63    pub follow_symlinks: bool,
64    pub single_filesystem: bool,
65    pub count_hardlinks: bool,
66    pub ignore_patterns: Vec<String>,
67    pub filter_patterns: Vec<String>,
68    pub max_file_size: Option<u64>,
69    pub batch: Option<PathBuf>,
70    pub nul_delim: bool,
71    pub threads: Threads,
72}
73
74pub fn crawl_and_process<O: Op, PM: PageMap + Send + Sync, D: DisplayMode<PM>>(
75    paths: &[PathBuf],
76    crawl_config: &CrawlConfig,
77    op: &O,
78    range: &FileRange,
79    stats: &Stats,
80    display: &D,
81    cancellation: &Cancellation,
82) -> crate::Result<Vec<O::Output>> {
83    cancellation.check()?;
84    tracing::info!("starting {} on {} path(s)", O::LABEL, paths.len());
85    validate_patterns(crawl_config)?;
86    let mut seen_inodes = HashSet::new();
87    let mut files = Vec::new();
88    let collection = collect_paths(
89        paths,
90        crawl_config,
91        &mut seen_inodes,
92        stats,
93        cancellation,
94        |path| {
95            files.push(path);
96            Ok(())
97        },
98    );
99    if let Err(error) = collection {
100        display.finish();
101        return Err(error);
102    }
103
104    let next = AtomicUsize::new(0);
105    let stop = AtomicBool::new(false);
106    let outputs = Mutex::new(Vec::with_capacity(files.len()));
107    let first_error = Mutex::new(None);
108    let workers = crawl_config.threads.effective().min(files.len());
109    let worker_result: crate::Result<()> = std::thread::scope(|scope| {
110        let mut handles = Vec::with_capacity(workers);
111        for _ in 0..workers {
112            handles.push(
113                std::thread::Builder::new()
114                    .spawn_scoped(scope, || {
115                        while !stop.load(Ordering::Relaxed) {
116                            let index = next.fetch_add(1, Ordering::Relaxed);
117                            let Some(path) = files.get(index) else {
118                                break;
119                            };
120                            match display.process_one::<O>(op, path, range, stats, cancellation) {
121                                Ok(Some(output)) => {
122                                    outputs.lock().unwrap().push((index, output));
123                                }
124                                Ok(None) => {}
125                                Err(error) => {
126                                    stop.store(true, Ordering::Relaxed);
127                                    let mut first = first_error.lock().unwrap();
128                                    if first.is_none() {
129                                        *first = Some(error);
130                                    }
131                                }
132                            }
133                        }
134                    })
135                    .map_err(|error| crate::Error::io("spawn file worker", error))?,
136            );
137        }
138        for handle in handles {
139            handle.join().map_err(|_| crate::Error::WorkerPanic)?;
140        }
141        Ok(())
142    });
143
144    display.finish();
145    worker_result?;
146    if let Some(error) = first_error.into_inner().unwrap() {
147        return Err(error);
148    }
149    cancellation.check()?;
150    op.finish()?;
151    tracing::info!(
152        "done: {} files, {} pages",
153        stats.total_files.load(Ordering::Relaxed),
154        stats.total_pages.load(Ordering::Relaxed),
155    );
156    let mut outputs = outputs.into_inner().unwrap();
157    outputs.sort_unstable_by_key(|(index, _)| *index);
158    Ok(outputs.into_iter().map(|(_, output)| output).collect())
159}
160
161pub fn validate_patterns(config: &CrawlConfig) -> crate::Result<()> {
162    build_overrides(Path::new("."), config).map(drop)
163}
164
165fn collect_paths(
166    paths: &[PathBuf],
167    crawl_config: &CrawlConfig,
168    seen_inodes: &mut HashSet<(u64, u64)>,
169    stats: &Stats,
170    cancellation: &Cancellation,
171    mut emit: impl FnMut(PathBuf) -> crate::Result<()>,
172) -> crate::Result<()> {
173    cancellation.check()?;
174    let mut all_paths: Vec<PathBuf> = paths.to_vec();
175
176    if let Some(batch_path) = &crawl_config.batch {
177        let batch_paths = read_batch_paths(batch_path, crawl_config.nul_delim)
178            .map_err(|error| crate::Error::io(batch_path.display().to_string(), error))?;
179        all_paths.extend(batch_paths);
180    }
181
182    let needs_meta = crawl_config.max_file_size.is_some()
183        || !crawl_config.count_hardlinks
184        || crawl_config.single_filesystem;
185
186    for path in &all_paths {
187        if cancellation.is_cancelled() {
188            break;
189        }
190        let metadata = fs_err::metadata(path)
191            .map_err(|error| crate::Error::io(path.display().to_string(), error))?;
192        if metadata.is_dir() {
193            tracing::info!("crawling directory {}", path.display());
194            stats.total_dirs.fetch_add(1, Ordering::Relaxed);
195            walk_dir_entries(
196                path,
197                crawl_config,
198                needs_meta,
199                seen_inodes,
200                stats,
201                cancellation,
202                &mut emit,
203            )?;
204        } else if metadata.is_file() {
205            if explicit_file_allowed(path, &metadata, crawl_config, needs_meta, seen_inodes) {
206                emit(path.clone())?;
207            }
208        } else {
209            return Err(crate::Error::io(
210                path.display().to_string(),
211                io::Error::new(io::ErrorKind::InvalidInput, "not a file or directory"),
212            ));
213        }
214    }
215    Ok(())
216}
217
218struct TraversalRoot {
219    physical: PathBuf,
220    logical: PathBuf,
221    device: Option<u64>,
222    overrides: Option<Arc<ignore::overrides::Override>>,
223    ignored: Option<Arc<ignore::overrides::Override>>,
224}
225
226impl TraversalRoot {
227    fn new(
228        physical: PathBuf,
229        logical: PathBuf,
230        device: Option<u64>,
231        config: &CrawlConfig,
232    ) -> crate::Result<Self> {
233        Ok(Self {
234            physical,
235            overrides: build_overrides(&logical, config)?.map(Arc::new),
236            ignored: build_ignore_overrides(&logical, config)?.map(Arc::new),
237            logical,
238            device,
239        })
240    }
241
242    fn logical_path(&self, physical: &Path) -> PathBuf {
243        physical.strip_prefix(&self.physical).map_or_else(
244            |_| physical.to_owned(),
245            |relative| self.logical.join(relative),
246        )
247    }
248}
249
250fn walk_dir_entries(
251    root: &Path,
252    config: &CrawlConfig,
253    needs_meta: bool,
254    seen_inodes: &mut HashSet<(u64, u64)>,
255    stats: &Stats,
256    cancellation: &Cancellation,
257    mut emit: impl FnMut(PathBuf) -> crate::Result<()>,
258) -> crate::Result<()> {
259    use std::os::unix::fs::MetadataExt as _;
260
261    let root_metadata = fs_err::metadata(root)
262        .map_err(|error| crate::Error::io(root.display().to_string(), error))?;
263    let root_device = config.single_filesystem.then(|| root_metadata.dev());
264    let physical_root = if config.follow_symlinks && root.is_symlink() {
265        fs_err::canonicalize(root)
266            .map_err(|error| crate::Error::io(root.display().to_string(), error))?
267    } else {
268        root.to_owned()
269    };
270    let mut pending = VecDeque::from([TraversalRoot::new(
271        physical_root.clone(),
272        root.to_owned(),
273        root_device,
274        config,
275    )?]);
276    let mut visited = HashSet::new();
277    if config.follow_symlinks {
278        visited.insert(fs_err::canonicalize(&physical_root).unwrap_or(physical_root));
279    }
280
281    while let Some(root) = pending.pop_front() {
282        if cancellation.is_cancelled() {
283            return Err(crate::Error::Cancelled);
284        }
285
286        let root = Arc::new(root);
287        let descend_root = Arc::clone(&root);
288        let descend_cancellation = cancellation.clone();
289        let walk_stop = Arc::new(AtomicBool::new(false));
290        let descend_stop = Arc::clone(&walk_stop);
291        let entries = dua_core::walk(
292            &root.physical,
293            config.threads.effective(),
294            Order::Completion,
295            move |entry| {
296                if descend_cancellation.is_cancelled() || descend_stop.load(Ordering::Relaxed) {
297                    return false;
298                }
299                if entry.depth == 0 {
300                    return true;
301                }
302                if descend_root.device.is_some_and(|device| {
303                    entry
304                        .metadata
305                        .as_ref()
306                        .map_or(true, |metadata| metadata.dev() != device)
307                }) {
308                    return false;
309                }
310                let path = descend_root.logical_path(&entry.path());
311                !descend_root
312                    .ignored
313                    .as_ref()
314                    .is_some_and(|overrides| overrides.matched(&path, true).is_ignore())
315            },
316        );
317
318        let mut walk_error = None;
319        for entry_result in entries {
320            if walk_error.is_some() {
321                continue;
322            }
323            if cancellation.is_cancelled() {
324                walk_stop.store(true, Ordering::Relaxed);
325                walk_error = Some(crate::Error::Cancelled);
326                continue;
327            }
328            let result = (|| -> crate::Result<()> {
329                let entry = entry_result
330                    .map_err(|error| crate::Error::io(root.logical.display().to_string(), error))?;
331                if entry.depth == 0 && entry.file_type.is_dir() {
332                    return Ok(());
333                }
334
335                let physical_path = entry.path();
336                let logical_path = root.logical_path(&physical_path);
337                if entry.file_type.is_symlink() {
338                    if !config.follow_symlinks {
339                        return Ok(());
340                    }
341                    let metadata = fs_err::metadata(&physical_path).map_err(|error| {
342                        crate::Error::io(logical_path.display().to_string(), error)
343                    })?;
344                    if root.device.is_some_and(|device| metadata.dev() != device) {
345                        return Ok(());
346                    }
347                    if metadata.is_dir() {
348                        stats.total_dirs.fetch_add(1, Ordering::Relaxed);
349                        let target = fs_err::canonicalize(&physical_path).map_err(|error| {
350                            crate::Error::io(logical_path.display().to_string(), error)
351                        })?;
352                        if visited.insert(target.clone()) {
353                            pending.push_back(TraversalRoot::new(
354                                target,
355                                logical_path,
356                                root.device,
357                                config,
358                            )?);
359                        }
360                    } else if metadata.is_file()
361                        && path_allowed(&logical_path, false, config, root.overrides.as_deref())
362                        && file_allowed(Some(&metadata), config, needs_meta, seen_inodes)
363                    {
364                        emit(logical_path)?;
365                    }
366                    return Ok(());
367                }
368
369                let metadata = match entry.metadata.as_ref() {
370                    Ok(metadata) => Some(metadata),
371                    Err(error) if needs_meta => {
372                        return Err(crate::Error::io(
373                            logical_path.display().to_string(),
374                            io::Error::new(error.kind(), error.to_string()),
375                        ));
376                    }
377                    Err(_) => None,
378                };
379
380                if entry.file_type.is_dir() {
381                    stats.total_dirs.fetch_add(1, Ordering::Relaxed);
382                    return Ok(());
383                }
384
385                if !entry.file_type.is_file()
386                    || !path_allowed(&logical_path, false, config, root.overrides.as_deref())
387                    || !file_allowed(metadata, config, needs_meta, seen_inodes)
388                {
389                    return Ok(());
390                }
391
392                emit(logical_path)?;
393                Ok(())
394            })();
395            if let Err(error) = result {
396                walk_stop.store(true, Ordering::Relaxed);
397                walk_error = Some(error);
398            }
399        }
400        if let Some(error) = walk_error {
401            return Err(error);
402        }
403    }
404    Ok(())
405}
406
407fn explicit_file_allowed(
408    path: &Path,
409    metadata: &std::fs::Metadata,
410    config: &CrawlConfig,
411    needs_meta: bool,
412    seen_inodes: &mut HashSet<(u64, u64)>,
413) -> bool {
414    let root = path
415        .parent()
416        .filter(|parent| !parent.as_os_str().is_empty())
417        .unwrap_or_else(|| Path::new("."));
418    let overrides =
419        build_overrides(root, config).expect("path patterns were validated before traversal");
420    if !path_allowed(path, false, config, overrides.as_ref()) {
421        return false;
422    }
423    file_allowed(Some(metadata), config, needs_meta, seen_inodes)
424}
425
426fn path_allowed(
427    path: &Path,
428    is_dir: bool,
429    config: &CrawlConfig,
430    overrides: Option<&ignore::overrides::Override>,
431) -> bool {
432    overrides.is_none_or(|overrides| {
433        let matched = overrides.matched(path, is_dir);
434        !matched.is_ignore() && (config.filter_patterns.is_empty() || matched.is_whitelist())
435    })
436}
437
438fn file_allowed(
439    metadata: Option<&std::fs::Metadata>,
440    config: &CrawlConfig,
441    needs_meta: bool,
442    seen_inodes: &mut HashSet<(u64, u64)>,
443) -> bool {
444    if needs_meta && metadata.is_none() {
445        return false;
446    }
447
448    if let Some(max_size) = config.max_file_size
449        && let Some(metadata) = metadata
450        && metadata.len() > max_size
451    {
452        return false;
453    }
454
455    if !config.count_hardlinks {
456        #[cfg(unix)]
457        {
458            use std::os::unix::fs::MetadataExt;
459            if let Some(metadata) = metadata
460                && metadata.nlink() > 1
461                && !seen_inodes.insert((metadata.dev(), metadata.ino()))
462            {
463                return false;
464            }
465        }
466    }
467
468    true
469}
470
471fn build_overrides(
472    root: &Path,
473    config: &CrawlConfig,
474) -> crate::Result<Option<ignore::overrides::Override>> {
475    if config.ignore_patterns.is_empty() && config.filter_patterns.is_empty() {
476        return Ok(None);
477    }
478
479    let mut overrides = ignore::overrides::OverrideBuilder::new(root);
480    for pattern in &config.ignore_patterns {
481        overrides.add(&format!("!{pattern}"))?;
482    }
483    for pattern in &config.filter_patterns {
484        overrides.add(pattern)?;
485    }
486    Ok(Some(overrides.build()?))
487}
488
489fn build_ignore_overrides(
490    root: &Path,
491    config: &CrawlConfig,
492) -> crate::Result<Option<ignore::overrides::Override>> {
493    if config.ignore_patterns.is_empty() {
494        return Ok(None);
495    }
496
497    let mut overrides = ignore::overrides::OverrideBuilder::new(root);
498    for pattern in &config.ignore_patterns {
499        overrides.add(&format!("!{pattern}"))?;
500    }
501    Ok(Some(overrides.build()?))
502}
503
504pub fn read_batch_paths(path: &Path, nul_delim: bool) -> io::Result<Vec<PathBuf>> {
505    use std::os::unix::ffi::OsStrExt;
506
507    let reader: Box<dyn BufRead> = if path == Path::new("-") {
508        Box::new(io::stdin().lock())
509    } else {
510        Box::new(io::BufReader::new(fs_err::File::open(path)?))
511    };
512
513    let delim = if nul_delim { b'\0' } else { b'\n' };
514    reader
515        .split(delim)
516        .filter_map(|r| match r {
517            Ok(buf) if !buf.is_empty() => {
518                Some(Ok(PathBuf::from(std::ffi::OsStr::from_bytes(&buf))))
519            }
520            Ok(_) => None,
521            Err(e) => Some(Err(e)),
522        })
523        .collect()
524}
525
526#[cfg(test)]
527mod tests {
528    use super::*;
529    use crate::Cancellation;
530    use crate::mode::Cli;
531    use crate::ops::{FileContext, Op, ResidencyEffect};
532    use std::sync::atomic::AtomicUsize;
533
534    #[test]
535    fn threads_parse_and_display() {
536        assert_eq!("0".parse::<Threads>().unwrap(), Threads::All);
537        assert_eq!("4".parse::<Threads>().unwrap().to_string(), "4");
538        assert_eq!(Threads::from(1).get(), 1);
539        assert!(Threads::All.get() > 0);
540        assert_eq!(
541            Threads::Exact(NonZeroU16::new(u16::MAX).unwrap()).effective(),
542            Threads::All.get()
543        );
544    }
545
546    #[test]
547    fn cancelled_collection_emits_no_paths() {
548        let file = tempfile::NamedTempFile::new().unwrap();
549        let config = CrawlConfig {
550            follow_symlinks: false,
551            single_filesystem: false,
552            count_hardlinks: true,
553            ignore_patterns: Vec::new(),
554            filter_patterns: Vec::new(),
555            max_file_size: None,
556            batch: None,
557            nul_delim: false,
558            threads: Threads::default(),
559        };
560        let mut seen_inodes = HashSet::new();
561        let stats = Stats::new();
562        let cancellation = Cancellation::new();
563        cancellation.cancel();
564        let mut emitted = Vec::new();
565
566        let result = collect_paths(
567            &[file.path().to_owned()],
568            &config,
569            &mut seen_inodes,
570            &stats,
571            &cancellation,
572            |path| {
573                emitted.push(path);
574                Ok(())
575            },
576        );
577
578        assert!(emitted.is_empty());
579        assert!(matches!(result, Err(crate::Error::Cancelled)));
580    }
581
582    #[test]
583    fn file_operations_use_requested_parallelism() {
584        struct ConcurrentOp {
585            active: AtomicUsize,
586            max_active: AtomicUsize,
587        }
588
589        impl Op for ConcurrentOp {
590            const LABEL: &'static str = "test";
591            const EFFECT: ResidencyEffect = ResidencyEffect::Preserve;
592            type Output = ();
593
594            fn execute<PM: PageMap + Sync>(&self, _ctx: &FileContext<'_, PM>) -> crate::Result<()> {
595                let active = self.active.fetch_add(1, Ordering::SeqCst) + 1;
596                self.max_active.fetch_max(active, Ordering::SeqCst);
597                std::thread::sleep(std::time::Duration::from_millis(20));
598                self.active.fetch_sub(1, Ordering::SeqCst);
599                Ok(())
600            }
601        }
602
603        let dir = tempfile::tempdir().unwrap();
604        for index in 0..4 {
605            fs_err::write(dir.path().join(index.to_string()), vec![0u8; 4096]).unwrap();
606        }
607        let config = CrawlConfig {
608            follow_symlinks: false,
609            single_filesystem: false,
610            count_hardlinks: true,
611            ignore_patterns: Vec::new(),
612            filter_patterns: Vec::new(),
613            max_file_size: None,
614            batch: None,
615            nul_delim: false,
616            threads: Threads::Exact(NonZeroU16::new(2).unwrap()),
617        };
618        let op = ConcurrentOp {
619            active: AtomicUsize::new(0),
620            max_active: AtomicUsize::new(0),
621        };
622
623        crawl_and_process::<_, crate::mincore::DefaultPageMap, _>(
624            &[dir.path().to_owned()],
625            &config,
626            &op,
627            &FileRange::full(),
628            &Stats::new(),
629            &Cli,
630            &Cancellation::new(),
631        )
632        .unwrap();
633
634        assert_eq!(op.max_active.load(Ordering::SeqCst), 2);
635    }
636
637    #[test]
638    fn cancellation_drains_wide_traversal() {
639        let dir = tempfile::tempdir().unwrap();
640        for index in 0..10_000 {
641            fs_err::write(dir.path().join(index.to_string()), []).unwrap();
642        }
643        let config = CrawlConfig {
644            follow_symlinks: false,
645            single_filesystem: false,
646            count_hardlinks: true,
647            ignore_patterns: Vec::new(),
648            filter_patterns: Vec::new(),
649            max_file_size: None,
650            batch: None,
651            nul_delim: false,
652            threads: Threads::Exact(NonZeroU16::new(4).unwrap()),
653        };
654        let cancellation = Cancellation::new();
655        let worker_cancellation = cancellation.clone();
656        let root = dir.path().to_owned();
657        let (started_sender, started_receiver) = std::sync::mpsc::channel();
658        let (resume_sender, resume_receiver) = std::sync::mpsc::channel();
659        let (result_sender, result_receiver) = std::sync::mpsc::channel();
660        std::thread::spawn(move || {
661            let mut seen = HashSet::new();
662            let mut started = false;
663            let result = walk_dir_entries(
664                &root,
665                &config,
666                false,
667                &mut seen,
668                &Stats::new(),
669                &worker_cancellation,
670                |_| {
671                    if !started {
672                        started = true;
673                        started_sender.send(()).unwrap();
674                        resume_receiver.recv().unwrap();
675                    }
676                    Ok(())
677                },
678            );
679            let _ = result_sender.send(result);
680        });
681
682        started_receiver
683            .recv_timeout(std::time::Duration::from_secs(5))
684            .expect("traversal did not start");
685        cancellation.cancel();
686        resume_sender.send(()).unwrap();
687        let result = result_receiver
688            .recv_timeout(std::time::Duration::from_secs(5))
689            .expect("cancelled traversal deadlocked");
690        assert!(matches!(result, Err(crate::Error::Cancelled)));
691    }
692}