Skip to main content

pagers_core/
crawl.rs

1use std::io::{self, BufRead};
2use std::path::{Path, PathBuf};
3use std::sync::atomic::Ordering;
4
5use ignore::WalkBuilder;
6
7use crate::mincore::PageMap;
8use crate::mode::DisplayMode;
9use crate::ops::{FileRange, Op, Stats};
10use crate::par::{InodeSet, SeenInodes as _};
11
12#[cfg(feature = "rayon")]
13pub use crate::par::Threads;
14
15#[derive(Debug, Clone, PartialEq)]
16#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
17pub struct CrawlConfig {
18    pub follow_symlinks: bool,
19    pub single_filesystem: bool,
20    pub count_hardlinks: bool,
21    pub ignore_patterns: Vec<String>,
22    pub filter_patterns: Vec<String>,
23    pub max_file_size: Option<u64>,
24    pub batch: Option<PathBuf>,
25    pub nul_delim: bool,
26    #[cfg(feature = "rayon")]
27    pub threads: Threads,
28}
29
30pub fn crawl_and_process<O: Op, PM: PageMap + Send + Sync, D: DisplayMode<PM>>(
31    paths: &[PathBuf],
32    crawl_config: &CrawlConfig,
33    op: &O,
34    range: &FileRange,
35    stats: &Stats,
36    display: &D,
37) -> crate::Result<Vec<O::Output>> {
38    tracing::info!("starting {} on {} path(s)", O::LABEL, paths.len());
39    let seen_inodes = InodeSet::default();
40
41    #[cfg(feature = "rayon")]
42    {
43        use rayon::prelude::*;
44
45        let pool = rayon::ThreadPoolBuilder::new()
46            .num_threads(crawl_config.threads.num_threads())
47            .build()?;
48
49        let buf = std::thread::available_parallelism().map_or(16, |n| n.get() * 4);
50        let (tx, rx) = std::sync::mpsc::sync_channel::<PathBuf>(buf);
51
52        let outputs = pool.install(|| {
53            rayon::scope(|s| {
54                s.spawn(|_| {
55                    collect_file_paths_streaming(paths, crawl_config, &seen_inodes, stats, tx);
56                });
57
58                rx.into_iter()
59                    .par_bridge()
60                    .filter_map(|path| display.process_one::<O>(op, &path, range, stats))
61                    .collect::<Vec<_>>()
62            })
63        });
64
65        display.finish();
66        op.finish()?;
67        tracing::info!(
68            "done: {} files, {} pages",
69            stats.total_files.load(Ordering::Relaxed),
70            stats.total_pages.load(Ordering::Relaxed),
71        );
72        Ok(outputs)
73    }
74
75    #[cfg(not(feature = "rayon"))]
76    {
77        let file_paths = collect_file_paths(paths, crawl_config, &seen_inodes, stats);
78        let outputs = file_paths
79            .iter()
80            .filter_map(|path| display.process_one::<O>(op, path, range, stats))
81            .collect();
82        display.finish();
83        op.finish()?;
84        tracing::info!(
85            "done: {} files, {} pages",
86            stats.total_files.load(Ordering::Relaxed),
87            stats.total_pages.load(Ordering::Relaxed),
88        );
89        Ok(outputs)
90    }
91}
92
93#[cfg(feature = "rayon")]
94fn collect_file_paths_streaming(
95    paths: &[PathBuf],
96    crawl_config: &CrawlConfig,
97    seen_inodes: &InodeSet,
98    stats: &Stats,
99    tx: std::sync::mpsc::SyncSender<PathBuf>,
100) {
101    let mut all_paths: Vec<PathBuf> = paths.to_vec();
102
103    if let Some(batch_path) = &crawl_config.batch {
104        match read_batch_paths(batch_path, crawl_config.nul_delim) {
105            Ok(batch_paths) => all_paths.extend(batch_paths),
106            Err(e) => tracing::warn!("batch file: {e}"),
107        }
108    }
109
110    let needs_meta = crawl_config.max_file_size.is_some() || !crawl_config.count_hardlinks;
111
112    for path in &all_paths {
113        if path.is_dir() {
114            tracing::info!("crawling directory {}", path.display());
115            stats.total_dirs.fetch_add(1, Ordering::Relaxed);
116            walk_dir_entries(path, crawl_config, needs_meta, seen_inodes, |p| {
117                let _ = tx.send(p);
118            });
119        } else if path.is_file() {
120            let _ = tx.send(path.clone());
121        } else {
122            tracing::warn!("skipping {}: not a file or directory", path.display());
123        }
124    }
125}
126
127#[cfg(not(feature = "rayon"))]
128fn collect_file_paths(
129    paths: &[PathBuf],
130    crawl_config: &CrawlConfig,
131    seen_inodes: &InodeSet,
132    stats: &Stats,
133) -> Vec<PathBuf> {
134    let mut all_paths: Vec<PathBuf> = paths.to_vec();
135
136    if let Some(batch_path) = &crawl_config.batch {
137        match read_batch_paths(batch_path, crawl_config.nul_delim) {
138            Ok(batch_paths) => all_paths.extend(batch_paths),
139            Err(e) => tracing::warn!("batch file: {e}"),
140        }
141    }
142
143    let needs_meta = crawl_config.max_file_size.is_some() || !crawl_config.count_hardlinks;
144    let mut file_paths = Vec::new();
145
146    for path in &all_paths {
147        if path.is_dir() {
148            tracing::info!("crawling directory {}", path.display());
149            stats.total_dirs.fetch_add(1, Ordering::Relaxed);
150            walk_dir_entries(path, crawl_config, needs_meta, seen_inodes, |p| {
151                file_paths.push(p);
152            });
153        } else if path.is_file() {
154            file_paths.push(path.clone());
155        } else {
156            tracing::warn!("skipping {}: not a file or directory", path.display());
157        }
158    }
159
160    tracing::info!("discovered {} files", file_paths.len());
161    file_paths
162}
163
164fn walk_dir_entries(
165    root: &Path,
166    config: &CrawlConfig,
167    needs_meta: bool,
168    seen_inodes: &InodeSet,
169    mut emit: impl FnMut(PathBuf),
170) {
171    let mut builder = WalkBuilder::new(root);
172    builder
173        .follow_links(config.follow_symlinks)
174        .same_file_system(config.single_filesystem)
175        .hidden(false)
176        .git_ignore(false)
177        .git_global(false)
178        .git_exclude(false);
179
180    if !config.ignore_patterns.is_empty() || !config.filter_patterns.is_empty() {
181        let mut overrides = ignore::overrides::OverrideBuilder::new(root);
182        for pat in &config.ignore_patterns {
183            let _ = overrides.add(&format!("!{pat}"));
184        }
185        for pat in &config.filter_patterns {
186            let _ = overrides.add(pat);
187        }
188        if let Ok(ov) = overrides.build() {
189            builder.overrides(ov);
190        }
191    }
192
193    for entry in builder.build() {
194        let Ok(entry) = entry.inspect_err(|e| tracing::warn!("{e}")) else {
195            continue;
196        };
197
198        let Some(ft) = entry.file_type() else {
199            continue;
200        };
201
202        if !ft.is_file() {
203            continue;
204        }
205
206        let entry_path = entry.path();
207        let meta = if needs_meta {
208            entry_path.metadata().ok()
209        } else {
210            None
211        };
212
213        if let Some(max_size) = config.max_file_size
214            && let Some(ref m) = meta
215            && m.len() > max_size
216        {
217            continue;
218        }
219
220        if !config.count_hardlinks {
221            #[cfg(unix)]
222            {
223                use std::os::unix::fs::MetadataExt;
224                if let Some(ref m) = meta
225                    && m.nlink() > 1
226                    && seen_inodes.already_seen((m.dev(), m.ino()))
227                {
228                    continue;
229                }
230            }
231        }
232
233        emit(entry_path.to_path_buf());
234    }
235}
236
237pub fn read_batch_paths(path: &Path, nul_delim: bool) -> io::Result<Vec<PathBuf>> {
238    use std::os::unix::ffi::OsStrExt;
239
240    let reader: Box<dyn BufRead> = if path == Path::new("-") {
241        Box::new(io::stdin().lock())
242    } else {
243        Box::new(io::BufReader::new(std::fs::File::open(path)?))
244    };
245
246    let delim = if nul_delim { b'\0' } else { b'\n' };
247    reader
248        .split(delim)
249        .filter_map(|r| match r {
250            Ok(buf) if !buf.is_empty() => {
251                Some(Ok(PathBuf::from(std::ffi::OsStr::from_bytes(&buf))))
252            }
253            Ok(_) => None,
254            Err(e) => Some(Err(e)),
255        })
256        .collect()
257}