use std::io::{self, BufRead};
use std::path::{Path, PathBuf};
use std::sync::atomic::Ordering;
use ignore::WalkBuilder;
use crate::mincore::PageMap;
use crate::mode::DisplayMode;
use crate::ops::{FileRange, Op, Stats};
use crate::par::{InodeSet, SeenInodes as _};
#[cfg(feature = "rayon")]
pub use crate::par::Threads;
#[derive(Debug, Clone, PartialEq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
pub struct CrawlConfig {
pub follow_symlinks: bool,
pub single_filesystem: bool,
pub count_hardlinks: bool,
pub ignore_patterns: Vec<String>,
pub filter_patterns: Vec<String>,
pub max_file_size: Option<u64>,
pub batch: Option<PathBuf>,
pub nul_delim: bool,
#[cfg(feature = "rayon")]
pub threads: Threads,
}
pub fn crawl_and_process<O: Op, PM: PageMap + Send + Sync, D: DisplayMode<PM>>(
paths: &[PathBuf],
crawl_config: &CrawlConfig,
op: &O,
range: &FileRange,
stats: &Stats,
display: &D,
) -> crate::Result<Vec<O::Output>> {
tracing::info!("starting {} on {} path(s)", O::LABEL, paths.len());
let seen_inodes = InodeSet::default();
#[cfg(feature = "rayon")]
{
use rayon::prelude::*;
let pool = rayon::ThreadPoolBuilder::new()
.num_threads(crawl_config.threads.num_threads())
.build()?;
let buf = std::thread::available_parallelism().map_or(16, |n| n.get() * 4);
let (tx, rx) = std::sync::mpsc::sync_channel::<PathBuf>(buf);
let outputs = pool.install(|| {
rayon::scope(|s| {
s.spawn({
let tx = tx;
move |_| {
collect_paths(paths, crawl_config, &seen_inodes, stats, |p| {
let _ = tx.send(p);
});
}
});
rx.into_iter()
.par_bridge()
.filter_map(|path| display.process_one::<O>(op, &path, range, stats))
.collect::<Vec<_>>()
})
});
display.finish();
op.finish()?;
tracing::info!(
"done: {} files, {} pages",
stats.total_files.load(Ordering::Relaxed),
stats.total_pages.load(Ordering::Relaxed),
);
Ok(outputs)
}
#[cfg(not(feature = "rayon"))]
{
let mut file_paths = Vec::new();
collect_paths(paths, crawl_config, &seen_inodes, stats, |p| {
file_paths.push(p);
});
tracing::info!("discovered {} files", file_paths.len());
let outputs = file_paths
.iter()
.filter_map(|path| display.process_one::<O>(op, path, range, stats))
.collect();
display.finish();
op.finish()?;
tracing::info!(
"done: {} files, {} pages",
stats.total_files.load(Ordering::Relaxed),
stats.total_pages.load(Ordering::Relaxed),
);
Ok(outputs)
}
}
fn collect_paths(
paths: &[PathBuf],
crawl_config: &CrawlConfig,
seen_inodes: &InodeSet,
stats: &Stats,
mut emit: impl FnMut(PathBuf),
) {
let mut all_paths: Vec<PathBuf> = paths.to_vec();
if let Some(batch_path) = &crawl_config.batch {
match read_batch_paths(batch_path, crawl_config.nul_delim) {
Ok(batch_paths) => all_paths.extend(batch_paths),
Err(e) => tracing::warn!("batch file: {e}"),
}
}
let needs_meta = crawl_config.max_file_size.is_some() || !crawl_config.count_hardlinks;
for path in &all_paths {
if path.is_dir() {
tracing::info!("crawling directory {}", path.display());
stats.total_dirs.fetch_add(1, Ordering::Relaxed);
walk_dir_entries(path, crawl_config, needs_meta, seen_inodes, &mut emit);
} else if path.is_file() {
emit(path.clone());
} else {
tracing::warn!("skipping {}: not a file or directory", path.display());
}
}
}
fn walk_dir_entries(
root: &Path,
config: &CrawlConfig,
needs_meta: bool,
seen_inodes: &InodeSet,
mut emit: impl FnMut(PathBuf),
) {
let mut builder = WalkBuilder::new(root);
builder
.follow_links(config.follow_symlinks)
.same_file_system(config.single_filesystem)
.hidden(false)
.git_ignore(false)
.git_global(false)
.git_exclude(false);
if !config.ignore_patterns.is_empty() || !config.filter_patterns.is_empty() {
let mut overrides = ignore::overrides::OverrideBuilder::new(root);
for pat in &config.ignore_patterns {
let _ = overrides.add(&format!("!{pat}"));
}
for pat in &config.filter_patterns {
let _ = overrides.add(pat);
}
if let Ok(ov) = overrides.build() {
builder.overrides(ov);
}
}
for entry in builder.build() {
let Ok(entry) = entry.inspect_err(|e| tracing::warn!("{e}")) else {
continue;
};
let Some(ft) = entry.file_type() else {
continue;
};
if !ft.is_file() {
continue;
}
let entry_path = entry.path();
let meta = if needs_meta {
fs_err::metadata(entry_path).ok()
} else {
None
};
if let Some(max_size) = config.max_file_size
&& let Some(ref m) = meta
&& m.len() > max_size
{
continue;
}
if !config.count_hardlinks {
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt;
if let Some(ref m) = meta
&& m.nlink() > 1
&& seen_inodes.already_seen((m.dev(), m.ino()))
{
continue;
}
}
}
emit(entry_path.to_path_buf());
}
}
pub fn read_batch_paths(path: &Path, nul_delim: bool) -> io::Result<Vec<PathBuf>> {
use std::os::unix::ffi::OsStrExt;
let reader: Box<dyn BufRead> = if path == Path::new("-") {
Box::new(io::stdin().lock())
} else {
Box::new(io::BufReader::new(fs_err::File::open(path)?))
};
let delim = if nul_delim { b'\0' } else { b'\n' };
reader
.split(delim)
.filter_map(|r| match r {
Ok(buf) if !buf.is_empty() => {
Some(Ok(PathBuf::from(std::ffi::OsStr::from_bytes(&buf))))
}
Ok(_) => None,
Err(e) => Some(Err(e)),
})
.collect()
}