34 Commits

Author SHA1 Message Date
sreedevk
9effdc9681 tui improvements 2025-07-11 22:46:20 +00:00
sreedevk
936c631623 added another test for processor 2025-07-11 14:24:21 +00:00
sreedevk
0d62361132 added processor test 2025-07-11 14:22:36 +00:00
sreedevk
b103d7d634 test finally passes 2025-07-11 01:26:59 +00:00
sreedevk
c42538a8c6 removed vfs 2025-07-10 01:19:14 +00:00
sreedevk
05437e0256 added a failing test - need to fix it 2025-07-09 21:49:44 +00:00
sreedevk
5942cf0b6c minor improvements 2025-07-08 20:20:55 +00:00
sreedevk
79c0c2be3d render filequeue 2025-07-07 13:33:24 +00:00
sreedevk
49b46f37f3 reorganization 2025-07-07 11:21:52 +00:00
sreedevk
a37d59308f application wrapper added 2025-07-07 10:52:45 +00:00
sreedevk
8a5b540e40 added application module 2025-07-07 01:49:36 +00:00
sreedevk
5ce0a2774a added server threadpool and message passing 2025-07-07 00:00:25 +00:00
sreedevk
8d32ac5151 fix: proc count decrement added 2025-07-06 13:16:25 +00:00
sreedevk
0c70e27835 ignore .bacon-locations 2025-07-06 12:45:31 +00:00
sreedevk
53282c50df implemented new multi threaded scanner 2025-07-06 12:45:06 +00:00
sreedev
b4bf5d4fb3 docs updated 2024-07-05 04:34:12 +00:00
Sreedev Kodichath
028b868ea9 version 0.2.2 (#59)
* replaced fxhash with gxhash
2024-07-04 23:23:01 -04:00
Sreedev Kodichath
a56d194ee3 Merge pull request #56 from sreedevk/feature/add-json-output
Version 0.2.1
2023-11-14 18:38:50 -05:00
sreedev
080cd791dc removed early return 2023-11-14 18:37:48 -05:00
sreedev
b16236763c updated readme + version number 2023-11-14 18:36:27 -05:00
sreedev
3e407f69c8 added json output for further processing using other tools 2023-11-14 18:33:40 -05:00
Sreedev Kodichath
7d386c9420 Merge pull request #54 from sreedevk/distribution-improvements
Improve Distribution Methods & Create Compiled Binaries for More Platforms
2023-08-10 19:54:29 -04:00
sreedev
0dc681d4e7 added release yml 2023-08-10 19:51:57 -04:00
sreedev
442cb4b519 added cargo dist options 2023-08-10 13:05:29 -04:00
sreedev
8463e72f2d removed debug information from distrbution & release profiles to reduce binary size 2023-08-10 12:59:07 -04:00
sreedev
05c95bb67a added roadmap to readme 2023-08-02 09:41:15 -04:00
Sreedev Kodichath
062d44acd9 Merge pull request #53 from sreedevk/v0.2.0
v0.2.0  Architecture Improvements
2023-07-17 18:17:57 -04:00
sreedev
ee6655de9b removed examples 2023-07-17 18:14:13 -04:00
sreedev
ffa5295598 clippy 2023-07-17 14:40:10 -04:00
sreedev
6be8596992 restored interactive mode 2023-07-17 14:34:42 -04:00
sreedev
a40f251e30 subtract overflow issue fixed 2023-07-17 14:19:57 -04:00
sreedev
02d05172da minsize issues fixed 2023-07-17 14:13:18 -04:00
sreedev
dcc709a666 added filetype filter 2023-07-17 14:06:57 -04:00
sreedev
e3d48ec505 v0.2.0 2023-07-17 13:59:58 -04:00
24 changed files with 2136 additions and 944 deletions

2
.cargo/config.toml Normal file
View File

@@ -0,0 +1,2 @@
[build]
rustflags = ["-C", "target-feature=+aes,+sse2"]

View File

@@ -1,12 +1,12 @@
# CI that:
#
# * checks for a Git Tag that looks like a release ("v1.2.0")
# * creates a Github Release™
# * builds binaries/packages with cargo-dist
# * uploads those packages to the Github Release™
# * checks for a Git Tag that looks like a release
# * creates a Github Release™ and fills in its text
# * builds artifacts with cargo-dist (executable-zips, installers)
# * uploads those artifacts to the Github Release™
#
# Note that the Github Release™ will be created before the packages,
# so there will be a few minutes where the release has no packages
# Note that the Github Release™ will be created before the artifacts,
# so there will be a few minutes where the release has no artifacts
# and then they will slowly trickle in, possibly failing. To make
# this more pleasant we mark the release as a "draft" until all
# artifacts have been successfully uploaded. This allows you to
@@ -17,114 +17,116 @@ name: Release
permissions:
contents: write
# This task will run whenever you push a git tag that looks like
# a version number. We just look for `v` followed by at least one number
# and then whatever. so `v1`, `v1.0.0`, and `v1.0.0-prerelease` all work.
# This task will run whenever you push a git tag that looks like a version
# like "v1", "v1.2.0", "v0.1.0-prerelease01", "my-app-v1.0.0", etc.
# The version will be roughly parsed as ({PACKAGE_NAME}-)?v{VERSION}, where
# PACKAGE_NAME must be the name of a Cargo package in your workspace, and VERSION
# must be a Cargo-style SemVer Version.
#
# If there's a prerelease-style suffix to the version then the Github Release™
# will be marked as a prerelease (handled by taiki-e/create-gh-release-action).
# If PACKAGE_NAME is specified, then we will create a Github Release™ for that
# package (erroring out if it doesn't have the given version or isn't cargo-dist-able).
#
# Note that when generating links to uploaded artifacts, cargo-dist will currently
# assume that your git tag is always v{VERSION} where VERSION is the version in
# the published package's Cargo.toml (this is the default behaviour of cargo-release).
# In the future this may be made more robust/configurable.
# If PACKAGE_NAME isn't specified, then we will create a Github Release™ for all
# (cargo-dist-able) packages in the workspace with that version (this is mode is
# intended for workspaces with only one dist-able package, or with all dist-able
# packages versioned/released in lockstep).
#
# If you push multiple tags at once, separate instances of this workflow will
# spin up, creating an independent Github Release™ for each one.
#
# If there's a prerelease-style suffix to the version then the Github Release™
# will be marked as a prerelease.
on:
push:
tags:
- v[0-9]+.*
env:
ALL_CARGO_DIST_TARGET_ARGS: --target=x86_64-unknown-linux-gnu --target=x86_64-apple-darwin --target=x86_64-pc-windows-msvc
ALL_CARGO_DIST_INSTALLER_ARGS:
- '*-?v[0-9]+*'
jobs:
# Create the Github Release™ so the packages have something to be uploaded to
# Create the Github Release™ so the packages have something to be uploaded to
create-release:
runs-on: ubuntu-latest
outputs:
tag: ${{ steps.create-gh-release.outputs.computed-prefix }}${{ steps.create-gh-release.outputs.version }}
has-releases: ${{ steps.create-release.outputs.has-releases }}
env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
steps:
- uses: actions/checkout@v3
- id: create-gh-release
uses: taiki-e/create-gh-release-action@v1
with:
draft: true
# (required) GitHub token for creating GitHub Releases.
token: ${{ secrets.GITHUB_TOKEN }}
- name: Install Rust
run: rustup update 1.71.0 --no-self-update && rustup default 1.71.0
- name: Install cargo-dist
run: curl --proto '=https' --tlsv1.2 -LsSf https://github.com/axodotdev/cargo-dist/releases/download/v0.0.7/cargo-dist-installer.sh | sh
- id: create-release
run: |
cargo dist plan --tag=${{ github.ref_name }} --output-format=json > dist-manifest.json
echo "dist plan ran successfully"
cat dist-manifest.json
# Create the Github Release™ based on what cargo-dist thinks it should be
ANNOUNCEMENT_TITLE=$(jq --raw-output ".announcement_title" dist-manifest.json)
IS_PRERELEASE=$(jq --raw-output ".announcement_is_prerelease" dist-manifest.json)
jq --raw-output ".announcement_github_body" dist-manifest.json > new_dist_announcement.md
gh release create ${{ github.ref_name }} --draft --prerelease="$IS_PRERELEASE" --title="$ANNOUNCEMENT_TITLE" --notes-file=new_dist_announcement.md
echo "created announcement!"
# Upload the manifest to the Github Release™
gh release upload ${{ github.ref_name }} dist-manifest.json
echo "uploaded manifest!"
# Disable all the upload-artifacts tasks if we have no actual releases
HAS_RELEASES=$(jq --raw-output ".releases != null" dist-manifest.json)
echo "has-releases=$HAS_RELEASES" >> "$GITHUB_OUTPUT"
# Build and packages all the things
upload-artifacts:
# Let the initial task tell us to not run (currently very blunt)
needs: create-release
if: ${{ needs.create-release.outputs.has-releases == 'true' }}
strategy:
matrix:
# For these target platforms
include:
- target: x86_64-unknown-linux-gnu
os: ubuntu-20.04
install-dist: curl --proto '=https' --tlsv1.2 -L -sSf https://github.com/axodotdev/cargo-dist/releases/download/v0.0.2/installer.sh | sh
- target: x86_64-apple-darwin
os: macos-11
install-dist: curl --proto '=https' --tlsv1.2 -L -sSf https://github.com/axodotdev/cargo-dist/releases/download/v0.0.2/installer.sh | sh
- target: x86_64-pc-windows-msvc
os: windows-2019
install-dist: irm 'https://github.com/axodotdev/cargo-dist/releases/download/v0.0.2/installer.ps1' | iex
- os: macos-11
dist-args: --artifacts=local --target=aarch64-apple-darwin --target=x86_64-apple-darwin
install-dist: curl --proto '=https' --tlsv1.2 -LsSf https://github.com/axodotdev/cargo-dist/releases/download/v0.0.7/cargo-dist-installer.sh | sh
- os: ubuntu-20.04
dist-args: --artifacts=local --target=x86_64-unknown-linux-gnu
install-dist: curl --proto '=https' --tlsv1.2 -LsSf https://github.com/axodotdev/cargo-dist/releases/download/v0.0.7/cargo-dist-installer.sh | sh
- os: windows-2019
dist-args: --artifacts=local --target=x86_64-pc-windows-msvc
install-dist: irm https://github.com/axodotdev/cargo-dist/releases/download/v0.0.7/cargo-dist-installer.ps1 | iex
runs-on: ${{ matrix.os }}
env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
steps:
- uses: actions/checkout@v3
- name: Install Rust
run: rustup update stable && rustup default stable
run: rustup update 1.71.0 --no-self-update && rustup default 1.71.0
- name: Install cargo-dist
run: ${{ matrix.install-dist }}
- name: Run cargo-dist
# This logic is a bit janky because it's trying to be a polyglot between
# powershell and bash since this will run on windows, macos, and linux!
# The two platforms don't agree on how to talk about env vars but they
# do agree on 'cat' and '$()' so we use that to marshal values between commmands.
# do agree on 'cat' and '$()' so we use that to marshal values between commands.
run: |
# Actually do builds and make zips and whatnot
cargo dist --target=${{ matrix.target }} --output-format=json > dist-manifest.json
cargo dist build --tag=${{ github.ref_name }} --output-format=json ${{ matrix.dist-args }} > dist-manifest.json
echo "dist ran successfully"
cat dist-manifest.json
# Parse out what we just built and upload it to the Github Release™
cat dist-manifest.json | jq --raw-output ".releases[].artifacts[].path" > uploads.txt
# Parse out what we just built and upload it to the Github Release™
jq --raw-output ".artifacts[]?.path | select( . != null )" dist-manifest.json > uploads.txt
echo "uploading..."
cat uploads.txt
gh release upload ${{ needs.create-release.outputs.tag }} $(cat uploads.txt)
gh release upload ${{ github.ref_name }} $(cat uploads.txt)
echo "uploaded!"
# Compute and upload the manifest for everything
upload-manifest:
needs: create-release
runs-on: ubuntu-latest
env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
steps:
- uses: actions/checkout@v3
- name: Install Rust
run: rustup update stable && rustup default stable
- name: Install cargo-dist
run: curl --proto '=https' --tlsv1.2 -L -sSf https://github.com/axodotdev/cargo-dist/releases/download/v0.0.2/installer.sh | sh
- name: Run cargo-dist manifest
run: |
# Generate a manifest describing everything
cargo dist manifest --no-local-paths --output-format=json $ALL_CARGO_DIST_TARGET_ARGS $ALL_CARGO_DIST_INSTALLER_ARGS > dist-manifest.json
echo "dist manifest ran successfully"
cat dist-manifest.json
# Upload the manifest to the Github Release™
gh release upload ${{ needs.create-release.outputs.tag }} dist-manifest.json
echo "uploaded manifest!"
# Edit the Github Release™ title/body to match what cargo-dist thinks it should be
CHANGELOG_TITLE=$(cat dist-manifest.json | jq --raw-output ".releases[].changelog_title")
cat dist-manifest.json | jq --raw-output ".releases[].changelog_body" > new_dist_changelog.md
gh release edit ${{ needs.create-release.outputs.tag }} --title="$CHANGELOG_TITLE" --notes-file=new_dist_changelog.md
echo "updated release notes!"
# Mark the Github Release™ as a non-draft now that everything has succeeded!
# Mark the Github Release™ as a non-draft now that everything has succeeded!
publish-release:
needs: [create-release, upload-artifacts, upload-manifest]
# Only run after all the other tasks, but it's ok if upload-artifacts was skipped
needs: [create-release, upload-artifacts]
if: ${{ always() && needs.create-release.result == 'success' && (needs.upload-artifacts.result == 'skipped' || needs.upload-artifacts.result == 'success') }}
runs-on: ubuntu-latest
env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
@@ -132,5 +134,4 @@ jobs:
- uses: actions/checkout@v3
- name: mark release as non-draft
run: |
gh release edit ${{ needs.create-release.outputs.tag }} --draft=false
gh release edit ${{ github.ref_name }} --draft=false

5
.gitignore vendored
View File

@@ -1,3 +1,4 @@
/target
/test_data
.envrc
/Cargo.lock
/.bacon-locations
/result-bin

1232
Cargo.lock generated

File diff suppressed because it is too large Load Diff

View File

@@ -1,8 +1,9 @@
[package]
name = "deduplicator"
version = "0.1.6"
version = "0.2.2"
edition = "2021"
description = "find,filter,delete Duplicates"
repository = "https://github.com/sreedevk/deduplicator"
license = "MIT"
authors = [
"Sreedev Kodichath <sreedevpadmakumar@gmail.com>",
@@ -19,18 +20,32 @@ chrono = "0.4.23"
clap = { version = "4.0.32", features = ["derive"] }
colored = "2.0.0"
dashmap = { version = "5.4.0", features = ["rayon"] }
fxhash = "0.2.1"
globwalk = "0.8.1"
gxhash = "3.4.1"
indicatif = { version = "0.17.2", features = ["rayon"] }
itertools = "0.10.5"
memmap2 = "0.5.8"
pathdiff = "0.2.1"
prettytable-rs = "0.10.0"
ratatui = "0.29.0"
rayon = "1.6.1"
serde = { version = "1.0.192", features = ["derive"] }
serde_json = "1.0.108"
tempfile = "3.20.0"
threadpool = "1.8.1"
unicode-segmentation = "1.10.0"
uuid = { version = "1.17.0", features = ["v4"] }
[profile.release]
strip = true
# generated by 'cargo dist init'
[profile.dist]
inherits = "release"
debug = true
split-debuginfo = "packed"
lto = "thin"
[workspace.metadata.dist]
rust-toolchain-version = "1.78.0"
ci = ["github"]
targets = ["x86_64-unknown-linux-gnu", "x86_64-apple-darwin", "x86_64-pc-windows-msvc", "aarch64-apple-darwin"]
cargo-dist-version = "0.0.7"

View File

@@ -21,6 +21,7 @@ Options:
-f, --follow-links Follow links while scanning directories
-h, --help Print help information
-V, --version Print version information
--json
```
### Examples
@@ -47,10 +48,15 @@ deduplicator ~/Media --min-size 100mb
#### Stable
> [!WARNING] Note from GxHash: GxHash relies on aes hardware acceleration, you must make sure the aes feature is enabled when building (otherwise it won't build). This can be done by setting the RUSTFLAGS environment variable to -C target-feature=+aes or -C target-cpu=native (the latter should work if your CPU is properly recognized by rustc, which is the case most of the time).
> please install version `0.2.1` if you are unable to install `0.2.2`
```bash
$ cargo install deduplicator
$ RUSTFLAGS="-C target-cpu=native" cargo install deduplicator
```
> [!]
#### Nightly
if you'd like to install with nightly features, you can use
@@ -121,3 +127,9 @@ Memory: 31731MiB (~32GiB)
## Screenshots
![](https://user-images.githubusercontent.com/36154121/213618143-e5182e39-731e-4817-87dd-1a6a0f38a449.gif)
## Roadmap
- Tree format output for duplicate file listing
- GUI
- Packages for different operating system repositories (currently only installable via cargo)
- TUI: Improve Key Handling

View File

@@ -1,17 +1,71 @@
use crate::output;
use crate::params::Params;
use crate::scanner;
use anyhow::Result;
use std::sync::mpsc::channel;
use std::sync::Arc;
pub struct App;
use anyhow::{anyhow, Result};
use threadpool::ThreadPool;
use crate::params::Params;
use crate::server::{Message, Server};
use crate::tui::Tui;
pub struct App {
tpool: ThreadPool,
server: Arc<Server>,
app_opts: Arc<Params>,
}
impl App {
pub fn init(app_args: &Params) -> Result<()> {
let duplicates = scanner::duplicates(app_args)?;
match app_args.interactive {
true => output::interactive(duplicates, app_args),
false => output::print(duplicates, app_args),
pub fn new(app_opts: Arc<Params>) -> Self {
Self {
tpool: ThreadPool::new(8),
server: Arc::new(Server::new().expect("server init failed")),
app_opts,
}
}
pub fn start(&self) -> Result<()> {
let (server_tx, server_rx) = channel::<Message>();
let (app_tx, app_rx) = channel::<Message>();
let server_ptr = self.server.clone();
let tui_server_ptr = self.server.clone();
let mut ui = Tui::new(app_tx.clone(), tui_server_ptr);
let root_dir = self
.app_opts
.get_directory()?
.into_os_string()
.into_string()
.map_err(|_| anyhow!("path to str conv failed"))?
.into_boxed_str();
self.tpool.execute(move || {
server_ptr.start(server_rx).expect("server init failed");
});
self.tpool.execute(move || {
ui.start().expect("ui init failed");
});
self.tpool.execute(move || loop {
match app_rx.try_recv() {
Ok(Message::Exit) => {
server_tx
.send(Message::Exit)
.expect("message passing to app from ui failed.");
break;
}
Ok(Message::AddScanDirectory(dir)) => {
server_tx
.send(Message::AddScanDirectory(dir))
.expect("message passing to server failed.");
}
_ => continue,
}
});
app_tx.send(Message::AddScanDirectory(root_dir))?;
self.tpool.join();
Ok(())
}

0
src/cli/mod.rs Normal file
View File

View File

@@ -1,21 +0,0 @@
use anyhow::Result;
use colored::Colorize;
use std::path::PathBuf;
#[derive(Debug, Clone)]
pub struct File {
pub path: PathBuf,
pub size: Option<u64>,
pub hash: Option<String>,
}
pub fn delete_files(files: Vec<File>) -> Result<()> {
files.into_iter().for_each(|file| {
match std::fs::remove_file(file.path.clone()) {
Ok(_) => println!("{}: {}", "DELETED".green(), file.path.display()),
Err(_) => println!("{}: {}", "FAILED".red(), file.path.display())
}
});
Ok(())
}

43
src/fileinfo.rs Normal file
View File

@@ -0,0 +1,43 @@
use anyhow::Result;
use gxhash::GxHasher;
use memmap2::Mmap;
use serde::Serialize;
use std::fs;
use std::hash::Hasher;
use std::{fs::Metadata, path::PathBuf};
#[derive(Debug, Clone, Serialize)]
pub struct FileInfo {
pub path: PathBuf,
pub hash: Option<String>,
pub size: u64,
#[serde(skip)]
pub filemeta: Metadata,
}
impl FileInfo {
pub fn hash(&self) -> Result<Self> {
let file = fs::File::open(self.path.clone())?;
let mapper = unsafe { Mmap::map(&file)? };
let mut primhasher = GxHasher::default();
mapper
.chunks(1_000_000)
.for_each(|chunk| primhasher.write(chunk));
Ok(Self {
hash: Some(primhasher.finish().to_string()),
..self.clone()
})
}
pub fn new(path: PathBuf) -> Result<Self> {
let filemeta = std::fs::metadata(path.clone())?;
Ok(Self {
path,
filemeta: filemeta.clone(),
hash: None,
size: filemeta.len(),
})
}
}

View File

@@ -1,12 +0,0 @@
use crate::file_manager::File;
use crate::params::Params;
pub fn is_file_gt_min_size(app_opts: &Params, file: &File) -> bool {
match app_opts.get_min_size() {
Some(msize) => match file.size {
Some(fsize) => fsize >= msize,
None => true,
},
None => true,
}
}

134
src/formatter.rs Normal file
View File

@@ -0,0 +1,134 @@
pub struct Formatter;
use crate::fileinfo::FileInfo;
use crate::params::Params;
use anyhow::Result;
use chrono::{DateTime, Utc};
use colored::Colorize;
use dashmap::DashMap;
use indicatif::{
ParallelProgressIterator, ProgressBar, ProgressFinish, ProgressIterator, ProgressStyle,
};
use pathdiff::diff_paths;
use prettytable::{format, row, Table};
use rayon::prelude::*;
use std::borrow::Cow;
use std::path::PathBuf;
use std::time::Duration;
impl Formatter {
pub fn human_path(
file: &FileInfo,
app_args: &Params,
min_path_length: usize,
) -> Result<String> {
let base_directory: PathBuf = app_args.get_directory()?;
let relative_path = diff_paths(file.path.clone(), base_directory).unwrap_or_default();
let formatted_path = format!(
"{:<0width$}",
relative_path.to_str().unwrap_or_default().to_string(),
width = min_path_length
);
Ok(formatted_path)
}
pub fn human_filesize(file: &FileInfo) -> Result<String> {
Ok(format!("{:>12}", bytesize::ByteSize::b(file.size)))
}
pub fn human_mtime(file: &FileInfo) -> Result<String> {
let modified_time: DateTime<Utc> = file.filemeta.modified()?.into();
Ok(modified_time.format("%Y-%m-%d %H:%M:%S").to_string())
}
pub fn generate_table(raw: Vec<FileInfo>, app_args: &Params) -> Result<Table> {
let basepath_length = app_args.get_directory()?.to_str().unwrap_or_default().len();
let max_filepath_length = raw
.iter()
.map(|file| file.path.to_str().unwrap_or_default().len())
.max()
.unwrap_or_default();
let min_path_length = if max_filepath_length > basepath_length {
max_filepath_length - basepath_length
} else {
0
};
let progress_style = ProgressStyle::with_template(
"[{elapsed_precise}] {bar:40.cyan/blue} {pos:>7}/{len:7} {msg}",
)?;
let progress_bar = ProgressBar::new(raw.len() as u64);
progress_bar.set_style(progress_style);
progress_bar.enable_steady_tick(Duration::from_millis(50));
progress_bar.set_message("reconciling data");
let duplicates_table: DashMap<String, Vec<FileInfo>> = DashMap::new();
raw.into_par_iter()
.progress_with(progress_bar)
.with_finish(ProgressFinish::WithMessage(Cow::from("data reconciled")))
.map(|file| file.hash())
.filter_map(Result::ok)
.for_each(|file| {
duplicates_table
.entry(file.hash.clone().unwrap_or_default())
.and_modify(|fileset| fileset.push(file.clone()))
.or_insert_with(|| vec![file]);
});
let mut output_table = Table::new();
output_table.set_titles(row!["hash", "duplicates"]);
let progress_style = ProgressStyle::with_template(
"[{elapsed_precise}] {bar:40.cyan/blue} {pos:>7}/{len:7} {msg}",
)?;
let progress_bar = ProgressBar::new(duplicates_table.len() as u64);
progress_bar.set_style(progress_style);
progress_bar.enable_steady_tick(Duration::from_millis(50));
progress_bar.set_message("generating output");
duplicates_table
.into_iter()
.progress_with(progress_bar)
.with_finish(ProgressFinish::WithMessage(Cow::from("output generated")))
.for_each(|(hash, group)| {
let mut inner_table = Table::new();
inner_table.set_format(*format::consts::FORMAT_NO_BORDER_LINE_SEPARATOR);
group.iter().for_each(|file| {
inner_table.add_row(row![
Self::human_path(file, app_args, min_path_length)
.unwrap_or_default()
.blue(),
Self::human_filesize(file).unwrap_or_default().red(),
Self::human_mtime(file).unwrap_or_default().yellow()
]);
});
output_table.add_row(row![hash.green(), inner_table]);
});
Ok(output_table)
}
pub fn print(raw: Vec<FileInfo>, app_args: &Params) -> Result<()> {
if raw.is_empty() {
println!(
"\n\n{}\n",
"No duplicates found matching your search criteria.".green()
);
return Ok(());
}
if app_args.json {
let output_json = serde_json::to_string_pretty(&raw)?;
println!("{}", output_json);
} else {
let output_table = Self::generate_table(raw, app_args)?;
output_table.printstd();
}
Ok(())
}
}

154
src/interactive.rs Normal file
View File

@@ -0,0 +1,154 @@
use crate::formatter::Formatter;
use crate::{fileinfo::FileInfo, params::Params};
use anyhow::Result;
use colored::Colorize;
use dashmap::DashMap;
use indicatif::{ParallelProgressIterator, ProgressBar, ProgressFinish, ProgressStyle};
use prettytable::{format, row, Table};
use rayon::prelude::*;
use std::{
borrow::Cow,
io::{self, Write},
time::Duration,
};
pub fn scan_group_confirmation() -> Result<bool> {
print!("\nconfirm? [y/N]: ");
std::io::stdout().flush()?;
let mut user_input = String::new();
io::stdin().read_line(&mut user_input)?;
match user_input.trim() {
"Y" | "y" => Ok(true),
_ => Ok(false),
}
}
pub fn scan_group_instruction() -> Result<String> {
println!("\nEnter the indices of the files you want to delete.");
println!("You can enter multiple files using commas to seperate file indices.");
println!("example: 1,2");
print!("\n> ");
std::io::stdout().flush()?;
let mut user_input = String::new();
io::stdin().read_line(&mut user_input)?;
Ok(user_input)
}
pub fn init(result: Vec<FileInfo>, app_args: &Params) -> Result<()> {
let basepath_length = app_args.get_directory()?.to_str().unwrap_or_default().len();
let max_filepath_length = result
.iter()
.map(|file| file.path.to_str().unwrap_or_default().len())
.max()
.unwrap_or_default();
let min_path_length = if max_filepath_length > basepath_length {
max_filepath_length - basepath_length
} else {
0
};
let progress_style = ProgressStyle::with_template(
"[{elapsed_precise}] {bar:40.cyan/blue} {pos:>7}/{len:7} {msg}",
)?;
let progress_bar = ProgressBar::new(result.len() as u64);
progress_bar.set_style(progress_style);
progress_bar.enable_steady_tick(Duration::from_millis(50));
progress_bar.set_message("reconciling data");
let duplicates: DashMap<String, Vec<FileInfo>> = DashMap::new();
result
.into_par_iter()
.progress_with(progress_bar)
.with_finish(ProgressFinish::WithMessage(Cow::from("data reconciled")))
.map(|file| file.hash())
.filter_map(Result::ok)
.for_each(|file| {
duplicates
.entry(file.hash.clone().unwrap_or_default())
.and_modify(|fileset| fileset.push(file.clone()))
.or_insert_with(|| vec![file]);
});
duplicates
.clone()
.into_iter()
.enumerate()
.for_each(|(gindex, (_, group))| {
let mut itable = Table::new();
itable.set_format(*format::consts::FORMAT_NO_BORDER_LINE_SEPARATOR);
itable.set_titles(row!["index", "filename", "size", "updated_at"]);
group.iter().enumerate().for_each(|(index, file)| {
itable.add_row(row![
index,
Formatter::human_path(file, app_args, min_path_length)
.unwrap_or_default()
.blue(),
Formatter::human_filesize(file).unwrap_or_default().red(),
Formatter::human_mtime(file).unwrap_or_default().yellow()
]);
});
process_group_action(&group, gindex, duplicates.len(), itable);
});
Ok(())
}
pub fn process_group_action(
duplicates: &Vec<FileInfo>,
dup_index: usize,
dup_size: usize,
table: Table,
) {
println!("\nDuplicate Set {} of {}\n", dup_index + 1, dup_size);
table.printstd();
let files_to_delete = scan_group_instruction().unwrap_or_default();
let parsed_file_indices = files_to_delete
.trim()
.split(',')
.filter(|element| !element.is_empty())
.map(|index| index.parse::<usize>().unwrap_or_default())
.collect::<Vec<usize>>();
if parsed_file_indices
.clone()
.into_iter()
.any(|index| index > (duplicates.len() - 1))
{
println!("{}", "Err: File Index Out of Bounds!".red());
return process_group_action(duplicates, dup_index, dup_size, table);
}
print!("{esc}[2J{esc}[1;1H", esc = 27 as char);
if parsed_file_indices.is_empty() {
return;
}
let files_to_delete = parsed_file_indices
.into_iter()
.map(|index| duplicates[index].clone());
println!("\n{}", "The following files will be deleted:".red());
files_to_delete
.clone()
.enumerate()
.for_each(|(index, file)| {
println!("{}: {}", index.to_string().blue(), file.path.display());
});
match scan_group_confirmation().unwrap() {
true => {
files_to_delete.into_iter().for_each(|file| {
match std::fs::remove_file(file.path.clone()) {
Ok(_) => println!("{}: {}", "DELETED".green(), file.path.display()),
Err(_) => println!("{}: {}", "FAILED".red(), file.path.display()),
}
});
}
false => println!("{}", "\nCancelled Delete Operation.".red()),
}
}

View File

@@ -1,14 +1,40 @@
mod app;
mod file_manager;
mod output;
mod fileinfo;
mod formatter;
mod interactive;
mod params;
mod processor;
mod scanner;
mod filters;
/* version 2.0 modules*/
mod cli;
mod server;
mod tui;
use anyhow::Result;
use app::App;
use self::app::App;
use clap::Parser;
use params::Params;
use std::sync::Arc;
// use formatter::Formatter;
// use processor::Processor;
// use scanner::Scanner;
fn main() -> Result<()> {
App::init(&params::Params::parse())
let app_args = Params::parse();
// let scan_results = Scanner::build(&app_args)?.scan()?;
// let processor = Processor::new(scan_results);
// let results = processor.sizewise()?.hashwise()?;
// match app_args.interactive {
// false => Formatter::print(results.files, &app_args)?,
// true => interactive::init(results.files, &app_args)?,
// }
App::new(Arc::new(app_args))
.start()
.expect("app init failed.");
Ok(())
}

View File

@@ -1,217 +0,0 @@
use crate::file_manager::{self, File};
use crate::params::Params;
use anyhow::Result;
use chrono::offset::Utc;
use chrono::DateTime;
use colored::Colorize;
use dashmap::DashMap;
use indicatif::{ProgressBar, ProgressIterator, ProgressStyle};
use itertools::Itertools;
use prettytable::{format, row, Table};
use std::io::Write;
use std::path::Path;
use std::time::Duration;
use std::{fs, io};
use unicode_segmentation::UnicodeSegmentation;
fn format_path(path: &Path, opts: &Params) -> Result<String> {
let display_path = path
.to_string_lossy()
.replace(opts.get_directory()?.to_string_lossy().as_ref(), "");
let display_range = if display_path.chars().count() > 32 {
display_path
.graphemes(true)
.collect::<Vec<&str>>()
.into_iter()
.rev()
.take(32)
.rev()
.collect()
} else {
display_path
};
Ok(format!("...{display_range:<32}"))
}
fn file_size(file: &File) -> Result<String> {
Ok(format!("{:>12}", bytesize::ByteSize::b(file.size.unwrap())))
}
fn modified_time(path: &Path) -> Result<String> {
let mdata = fs::metadata(path)?;
let modified_time: DateTime<Utc> = mdata.modified()?.into();
Ok(modified_time.format("%Y-%m-%d %H:%M:%S").to_string())
}
fn scan_group_instruction() -> Result<String> {
println!("\nEnter the indices of the files you want to delete.");
println!("You can enter multiple files using commas to seperate file indices.");
println!("example: 1,2");
print!("\n> ");
std::io::stdout().flush()?;
let mut user_input = String::new();
io::stdin().read_line(&mut user_input)?;
Ok(user_input)
}
fn scan_group_confirmation() -> Result<bool> {
print!("\nconfirm? [y/N]: ");
std::io::stdout().flush()?;
let mut user_input = String::new();
io::stdin().read_line(&mut user_input)?;
match user_input.trim() {
"Y" | "y" => Ok(true),
_ => Ok(false),
}
}
fn process_group_action(duplicates: &Vec<File>, dup_index: usize, dup_size: usize, table: Table) {
println!("\nDuplicate Set {} of {}\n", dup_index + 1, dup_size);
table.printstd();
let files_to_delete = scan_group_instruction().unwrap_or_default();
let parsed_file_indices = files_to_delete
.trim()
.split(',')
.filter(|element| !element.is_empty())
.map(|index| index.parse::<usize>().unwrap_or_default())
.collect::<Vec<usize>>();
if parsed_file_indices
.clone()
.into_iter()
.any(|index| index > (duplicates.len() - 1))
{
println!("{}", "Err: File Index Out of Bounds!".red());
return process_group_action(duplicates, dup_index, dup_size, table);
}
print!("{esc}[2J{esc}[1;1H", esc = 27 as char);
if parsed_file_indices.is_empty() {
return;
}
let files_to_delete = parsed_file_indices
.into_iter()
.map(|index| duplicates[index].clone());
println!("\n{}", "The following files will be deleted:".red());
files_to_delete
.clone()
.enumerate()
.for_each(|(index, file)| {
println!("{}: {}", index.to_string().blue(), file.path.display());
});
match scan_group_confirmation().unwrap() {
true => {
file_manager::delete_files(files_to_delete.collect::<Vec<File>>()).ok();
}
false => println!("{}", "\nCancelled Delete Operation.".red()),
}
}
pub fn interactive(duplicates: DashMap<String, Vec<File>>, opts: &Params) {
if duplicates.is_empty() {
println!(
"\n{}",
"No duplicates found matching your search criteria.".green()
);
return;
}
duplicates
.clone()
.into_iter()
.sorted_unstable_by_key(|(_, f)| {
-(f.first().and_then(|ff| ff.size).unwrap_or_default() as i64)
}) // sort by descending file size in interactive mode
.enumerate()
.for_each(|(gindex, (_, group))| {
let mut itable = Table::new();
itable.set_format(*format::consts::FORMAT_NO_BORDER_LINE_SEPARATOR);
itable.set_titles(row!["index", "filename", "size", "updated_at"]);
group.iter().enumerate().for_each(|(index, file)| {
itable.add_row(row![
index,
format_path(&file.path, opts).unwrap_or_default().blue(),
file_size(file).unwrap_or_default().red(),
modified_time(&file.path).unwrap_or_default().yellow()
]);
});
process_group_action(&group, gindex, duplicates.len(), itable);
});
}
pub fn print(duplicates: DashMap<String, Vec<File>>, opts: &Params) {
if duplicates.is_empty() {
println!(
"\n{}",
"No duplicates found matching your search criteria.".green()
);
return;
}
let mut output_table = Table::new();
let progress_bar = ProgressBar::new(duplicates.len() as u64);
progress_bar.enable_steady_tick(Duration::from_millis(50));
let progress_style = ProgressStyle::default_bar()
.template("{spinner:.green} [generating output] [{wide_bar:.cyan/blue}] {pos}/{len} files")
.unwrap();
progress_bar.set_style(progress_style);
output_table.set_titles(row!["hash", "duplicates"]);
duplicates
.into_iter()
.sorted_unstable_by_key(|(_, f)| f.first().and_then(|ff| ff.size).unwrap_or_default())
.progress_with(progress_bar)
.for_each(|(hash, group)| {
let mut inner_table = Table::new();
inner_table.set_format(*format::consts::FORMAT_NO_BORDER_LINE_SEPARATOR);
group.iter().for_each(|file| {
inner_table.add_row(row![
format_path(&file.path, opts).unwrap_or_default().blue(),
file_size(file).unwrap_or_default().red(),
modified_time(&file.path).unwrap_or_default().yellow()
]);
});
output_table.add_row(row![hash.green(), inner_table]);
});
output_table.printstd();
}
#[allow(unused)]
pub fn raw(duplicates: DashMap<String, Vec<File>>, opts: &Params) -> Result<()> {
if duplicates.is_empty() {
println!(
"\n{}",
"No duplicates found matching your search criteria.".green()
);
return Ok(());
}
duplicates
.into_iter()
.sorted_unstable_by_key(|(_, f)| f.first().and_then(|ff| ff.size).unwrap_or_default())
.for_each(|(_hash, group)| {
group.iter().for_each(|file| {
println!(
"{}\t{}\t{}",
format_path(&file.path, opts).unwrap_or_default().blue(),
file_size(file).unwrap_or_default().red(),
modified_time(&file.path).unwrap_or_default().yellow()
)
});
println!("---");
});
Ok(())
}

View File

@@ -1,10 +1,9 @@
use std::{fs, path::PathBuf};
use anyhow::{anyhow, Result};
use anyhow::Result;
use clap::{Parser, ValueHint};
use globwalk::{GlobWalker, GlobWalkerBuilder};
#[derive(Parser, Debug)]
#[derive(Parser, Debug, Clone)]
#[command(author, version, about, long_about = None)]
pub struct Params {
/// Filetypes to deduplicate [default = all]
@@ -28,6 +27,9 @@ pub struct Params {
/// Follow links while scanning directories
#[arg(long, short)]
pub follow_links: bool,
/// print json output
#[arg(long)]
pub json: bool,
}
impl Params {
@@ -48,41 +50,7 @@ impl Params {
Ok(dir)
}
fn add_glob_min_depth(&self, builder: GlobWalkerBuilder) -> Result<GlobWalkerBuilder> {
match self.min_depth {
Some(mindepth) => Ok(builder.min_depth(mindepth)),
None => Ok(builder),
}
}
fn add_glob_max_depth(&self, builder: GlobWalkerBuilder) -> Result<GlobWalkerBuilder> {
match self.max_depth {
Some(maxdepth) => Ok(builder.max_depth(maxdepth)),
None => Ok(builder),
}
}
fn add_glob_follow_links(&self, builder: GlobWalkerBuilder) -> Result<GlobWalkerBuilder> {
match self.follow_links {
true => Ok(builder.follow_links(true)),
false => Ok(builder.follow_links(false)),
}
}
pub fn get_glob_walker(&self) -> Result<GlobWalker> {
let pattern: String = match self.types.as_ref() {
Some(filetypes) => format!("**/*{{{filetypes}}}"),
None => "**/*".to_string(),
};
let glob_walker_builder = self
.add_glob_min_depth(GlobWalkerBuilder::from_patterns(
self.get_directory()?,
&[pattern],
))
.and_then(|builder| self.add_glob_max_depth(builder))
.and_then(|builder| self.add_glob_follow_links(builder))?;
glob_walker_builder.build().map_err(|e| anyhow!(e))
pub fn get_types(&self) -> Option<String> {
self.types.clone()
}
}

107
src/processor.rs Normal file
View File

@@ -0,0 +1,107 @@
use anyhow::Result;
use dashmap::DashMap;
use indicatif::{ParallelProgressIterator, ProgressBar, ProgressStyle, ProgressFinish};
use rayon::prelude::{IntoParallelIterator, ParallelIterator};
use std::{time::Duration, borrow::Cow};
use crate::fileinfo::FileInfo;
#[derive(Debug, Clone)]
pub enum State {
Initial,
SizeWise,
HashWise,
}
#[derive(Debug, Clone)]
pub struct Processor {
pub files: Vec<FileInfo>,
pub state: State,
}
impl Processor {
pub fn new(files: Vec<FileInfo>) -> Self {
Self {
files,
state: State::Initial,
}
}
pub fn hashwise(&self) -> Result<Self> {
if self.files.is_empty() {
return Ok(self.clone());
}
let progress_style = ProgressStyle::with_template("[{elapsed_precise}] {bar:40.cyan/blue} {pos:>7}/{len:7} {msg}")?;
let progress_bar = ProgressBar::new(self.files.len() as u64);
progress_bar.set_style(progress_style);
progress_bar.enable_steady_tick(Duration::from_millis(50));
progress_bar.set_message("indexing file hashes");
let duplicates_table: DashMap<String, Vec<FileInfo>> = DashMap::new();
self.files
.clone()
.into_par_iter()
.progress_with(progress_bar)
.with_finish(ProgressFinish::WithMessage(Cow::from("indexed files hashes")))
.map(|file| file.hash())
.filter_map(Result::ok)
.for_each(|file| {
duplicates_table
.entry(file.hash.clone().unwrap_or_default())
.and_modify(|fileset| fileset.push(file.clone()))
.or_insert_with(|| vec![file]);
});
let files = duplicates_table
.into_read_only()
.values()
.cloned()
.filter(|subfiles| subfiles.len() > 1)
.flatten()
.collect::<Vec<FileInfo>>();
Ok(Self {
files,
state: State::HashWise,
})
}
pub fn sizewise(&self) -> Result<Self> {
if self.files.is_empty() {
return Ok(self.clone());
}
let progress_style = ProgressStyle::with_template("[{elapsed_precise}] {bar:40.cyan/blue} {pos:>7}/{len:7} {msg}")?;
let progress_bar = ProgressBar::new(self.files.len() as u64);
progress_bar.set_style(progress_style);
progress_bar.enable_steady_tick(Duration::from_millis(50));
progress_bar.set_message("indexing file sizes");
let duplicates_table: DashMap<u64, Vec<FileInfo>> = DashMap::new();
self.files
.clone()
.into_par_iter()
.progress_with(progress_bar)
.with_finish(ProgressFinish::WithMessage(Cow::from("indexed files sizes")))
.for_each(|file| {
duplicates_table
.entry(file.size)
.and_modify(|fileset| fileset.push(file.clone()))
.or_insert_with(|| vec![file]);
});
let files = duplicates_table
.into_read_only()
.values()
.cloned()
.filter(|subfiles| subfiles.len() > 1)
.flatten()
.collect::<Vec<FileInfo>>();
Ok(Self {
files,
state: State::SizeWise,
})
}
}

View File

@@ -1,151 +1,173 @@
use crate::{file_manager::File, filters, params::Params};
#![allow(unused)]
use crate::{fileinfo::FileInfo, params::Params};
use anyhow::Result;
use dashmap::DashMap;
use fxhash::hash64 as hasher;
use indicatif::{ParallelProgressIterator, ProgressBar, ProgressIterator, ProgressStyle};
use memmap2::Mmap;
use rayon::prelude::*;
use std::hash::Hasher;
use std::time::Duration;
use std::{
fs,
path::{Path, PathBuf},
};
use indicatif::{ProgressBar, ProgressStyle};
use std::{fs, path::PathBuf, time::Duration};
#[derive(Clone, Copy)]
enum IndexCritera {
Size,
Hash,
use globwalk::{GlobWalker, GlobWalkerBuilder};
#[derive(Debug, Clone)]
pub struct Scanner {
pub directory: Option<PathBuf>,
pub filetypes: Option<String>,
pub min_depth: Option<usize>,
pub max_depth: Option<usize>,
pub min_size: Option<u64>,
pub follow_links: bool,
}
pub fn duplicates(app_opts: &Params) -> Result<DashMap<String, Vec<File>>> {
let scan_results = scan(app_opts)?;
let size_index_store = index_files(scan_results, IndexCritera::Size)?;
let sizewize_duplicate_files = size_index_store
.into_par_iter()
.filter(|(_, files)| files.len() > 1)
.map(|(_, files)| files)
.flatten()
.collect::<Vec<File>>();
if sizewize_duplicate_files.len() > 1 {
let hash_index_store = index_files(sizewize_duplicate_files, IndexCritera::Hash)?;
let duplicate_files = hash_index_store
.into_par_iter()
.filter(|(_, files)| files.len() > 1)
.collect();
Ok(duplicate_files)
} else {
Ok(DashMap::new())
impl Scanner {
pub fn new() -> Self {
Self {
directory: None,
filetypes: None,
min_depth: None,
max_depth: None,
min_size: None,
follow_links: true,
}
}
}
fn scan(app_opts: &Params) -> Result<Vec<File>> {
let walker = app_opts.get_glob_walker()?;
let progress = ProgressBar::new_spinner();
let progress_style =
ProgressStyle::with_template("{spinner:.green} [mapping paths] {pos} paths")?;
progress.set_style(progress_style);
progress.enable_steady_tick(Duration::from_millis(50));
pub fn build(app_args: &Params) -> Result<Self> {
let scan_directory = app_args.get_directory()?;
Ok(Scanner::new())
.map(|scanner| scanner.directory(scan_directory))
.map(|scanner| match app_args.get_min_size() {
Some(min_size) => scanner.min_size(min_size),
None => scanner,
})
.map(|scanner| match app_args.get_types() {
Some(ftypes) => scanner.filetypes(ftypes),
None => scanner,
})
.map(|scanner| match app_args.min_depth {
Some(min_depth) => scanner.min_depth(min_depth),
None => scanner,
})
.map(|scanner| match app_args.max_depth {
Some(max_depth) => scanner.max_depth(max_depth),
None => scanner,
})
}
let files = walker
.progress_with(progress)
.filter_map(Result::ok)
.map(|file| file.into_path())
.filter(|fpath| fpath.is_file())
.collect::<Vec<PathBuf>>();
pub fn min_size(&self, min_size: u64) -> Self {
Self {
min_size: Some(min_size),
..self.clone()
}
}
let scan_progress = ProgressBar::new(files.len() as u64);
let scan_progress_style = ProgressStyle::with_template(
"{spinner:.green} [processing mapped paths] [{wide_bar:.cyan/blue}] {pos}/{len} files",
)?;
scan_progress.set_style(scan_progress_style);
scan_progress.enable_steady_tick(Duration::from_millis(50));
pub fn min_depth(&self, min_depth: usize) -> Self {
Self {
min_depth: Some(min_depth),
..self.clone()
}
}
let scan_results = files
.into_par_iter()
.progress_with(scan_progress)
.map(|fpath| File {
path: fpath.clone(),
hash: None,
size: Some(
fs::metadata(fpath)
.map(|metadata| metadata.len())
.unwrap_or_default(),
),
pub fn max_depth(&self, max_depth: usize) -> Self {
Self {
max_depth: Some(max_depth),
..self.clone()
}
}
pub fn directory(&self, dir: PathBuf) -> Self {
Self {
directory: Some(dir),
..self.clone()
}
}
pub fn filetypes(&self, patterns: String) -> Self {
Self {
filetypes: Some(patterns),
..self.clone()
}
}
pub fn ignore_links(&self) -> Self {
Self {
follow_links: false,
..self.clone()
}
}
pub fn follow_links(&self) -> Self {
Self {
follow_links: true,
..self.clone()
}
}
fn scan_patterns(&self) -> Result<String> {
Ok(match self.filetypes.clone() {
Some(ftypes) => format!("**/*{{{ftypes}}}"),
None => "**/*".to_string(),
})
.filter(|file| filters::is_file_gt_min_size(app_opts, file))
.collect();
}
Ok(scan_results)
}
fn scan_dir(&self) -> Result<PathBuf> {
let scan_dir = match self.directory.clone() {
Some(path) => path,
None => std::env::current_dir()?,
};
fn process_file_index(
mut file: File,
store: &DashMap<String, Vec<File>>,
index_criteria: IndexCritera,
) {
match index_criteria {
IndexCritera::Size => {
store
.entry(file.size.unwrap_or_default().to_string())
.and_modify(|fileset| fileset.push(file.clone()))
.or_insert_with(|| vec![file]);
}
IndexCritera::Hash => {
file.hash = Some(hash_file(&file.path).unwrap_or_default());
store
.entry(file.clone().hash.unwrap_or_default())
.and_modify(|fileset| fileset.push(file.clone()))
.or_insert_with(|| vec![file]);
Ok(fs::canonicalize(scan_dir)?)
}
fn attach_link_opts(&self, walker: GlobWalkerBuilder) -> Result<GlobWalkerBuilder> {
Ok(walker.follow_links(self.follow_links))
}
fn attach_walker_min_depth(&self, walker: GlobWalkerBuilder) -> Result<GlobWalkerBuilder> {
match self.min_depth {
Some(min_depth) => Ok(walker.min_depth(min_depth)),
None => Ok(walker),
}
}
}
fn index_files(
files: Vec<File>,
index_criteria: IndexCritera,
) -> Result<DashMap<String, Vec<File>>> {
let store: DashMap<String, Vec<File>> = DashMap::new();
let index_progress = ProgressBar::new(files.len() as u64);
let index_progress_style = ProgressStyle::with_template(
"{spinner:.green} [indexing files] [{wide_bar:.cyan/blue}] {pos}/{len} files",
)?;
index_progress.set_style(index_progress_style);
index_progress.enable_steady_tick(Duration::from_millis(50));
fn attach_walker_max_depth(&self, walker: GlobWalkerBuilder) -> Result<GlobWalkerBuilder> {
match self.max_depth {
Some(max_depth) => Ok(walker.max_depth(max_depth)),
None => Ok(walker),
}
}
fn build_walker(&self) -> Result<GlobWalker> {
let walker = Ok(GlobWalkerBuilder::from_patterns(
self.scan_dir()?,
&[self.scan_patterns()?],
))
.and_then(|walker| self.attach_walker_min_depth(walker))
.and_then(|walker| self.attach_walker_max_depth(walker))
.and_then(|walker| self.attach_link_opts(walker))?;
files
.into_par_iter()
.progress_with(index_progress)
.for_each(|file| process_file_index(file, &store, index_criteria));
Ok(walker.build()?)
}
Ok(store)
}
pub fn scan(&self) -> Result<Vec<FileInfo>> {
let progress_style = ProgressStyle::with_template("[{elapsed_precise}] {pos:>7} {msg}")?;
let progress_bar = ProgressBar::new_spinner();
progress_bar.set_style(progress_style);
progress_bar.enable_steady_tick(Duration::from_millis(50));
progress_bar.set_message("paths mapped");
let min_size = self.min_size.unwrap_or_default();
fn incremental_hashing(filepath: &Path) -> Result<String> {
let file = fs::File::open(filepath)?;
let fmap = unsafe { Mmap::map(&file)? };
let mut inchasher = fxhash::FxHasher::default();
let results = self
.build_walker()?
.filter_map(Result::ok)
.map(|entity| entity.into_path())
.map(|path| {
progress_bar.inc(1);
path
})
.filter(|path| path.is_file())
.map(FileInfo::new)
.filter_map(Result::ok)
.filter(|file| file.size > min_size)
.collect::<Vec<FileInfo>>();
fmap.chunks(1_000_000)
.for_each(|mega| inchasher.write(mega));
progress_bar.finish_with_message("paths mapped");
Ok(format!("{}", inchasher.finish()))
}
fn standard_hashing(filepath: &Path) -> Result<String> {
let file = fs::read(filepath)?;
Ok(hasher(&*file).to_string())
}
fn hash_file(filepath: &Path) -> Result<String> {
let filemeta = fs::metadata(filepath)?;
// NOTE: USE INCREMENTAL HASHING ONLY FOR FILES > 100MB
match filemeta.len() < 100_000_000 {
true => standard_hashing(filepath),
false => incremental_hashing(filepath),
Ok(results)
}
}

44
src/server/file.rs Normal file
View File

@@ -0,0 +1,44 @@
use anyhow::Result;
use std::fs::File;
use std::io::Read;
use std::os::unix::fs::MetadataExt;
use std::sync::Arc;
use uuid::Uuid;
const PARTIAL_SIZE: u64 = 4096;
#[allow(unused)]
#[derive(Clone, Debug)]
pub struct FileMeta {
pub id: Uuid,
pub path: Box<str>,
pub size: u64,
pub modtime: i64,
pub partial: Arc<[u8]>,
pub full_hash: Arc<[u8]>,
}
impl FileMeta {
fn create_partial(path: &str) -> Result<Arc<[u8]>> {
let mut partial_take = File::open(path)?.take(PARTIAL_SIZE);
let mut partial_buffer = Vec::with_capacity(PARTIAL_SIZE as usize);
partial_take.read_to_end(&mut partial_buffer)?;
Ok(Arc::from(partial_buffer))
}
pub fn new(path: Box<str>) -> Result<Self> {
let pstr = path.to_string();
let filemeta = std::fs::metadata(&pstr)?;
let partial = Self::create_partial(&pstr)?;
Ok(Self {
id: Uuid::new_v4(),
size: filemeta.size(),
modtime: filemeta.mtime(),
full_hash: Arc::new([]),
path,
partial,
})
}
}

79
src/server/mod.rs Normal file
View File

@@ -0,0 +1,79 @@
pub mod file;
mod processor;
mod scanner;
mod store;
use anyhow::Result;
use std::sync::mpsc;
use std::sync::mpsc::channel;
use std::sync::Arc;
use std::sync::Mutex;
use threadpool::ThreadPool;
use self::processor::Processor;
use self::scanner::Scanner;
use self::store::Store;
pub type FileQueue = Arc<Mutex<Vec<Box<str>>>>;
pub enum Message {
AddScanDirectory(Box<str>),
Exit,
None,
}
pub struct Server {
pub fq: FileQueue,
pub dupstore: Arc<Store>,
pub tpool: ThreadPool,
}
impl Server {
pub fn new() -> Result<Self> {
Ok(Self {
fq: Arc::new(Mutex::new(vec![])),
dupstore: Arc::new(Store::new()),
tpool: ThreadPool::new(4),
})
}
pub fn start(&self, rx: mpsc::Receiver<Message>) -> Result<()> {
let processor_fq = self.fq.clone();
let processor_store = self.dupstore.clone();
let (processor_tx, processor_rx) = channel::<Message>();
self.tpool.execute(move || {
Processor::new(processor_fq, processor_store, processor_rx)
.process()
.expect("processer execution interrupted.");
});
let scanner_fq = self.fq.clone();
let (scanner_tx, scanner_rx) = channel::<Message>();
self.tpool.execute(move || {
Scanner::new(scanner_fq, scanner_rx)
.index()
.expect("scanner indexing interrupted.");
});
self.tpool.execute(move || loop {
match rx.recv() {
Ok(Message::AddScanDirectory(path)) => {
scanner_tx
.send(Message::AddScanDirectory(path))
.expect("scanner tx message passing failed.");
}
Ok(Message::None) => {}
Ok(Message::Exit) | Err(_) => {
scanner_tx.send(Message::Exit).unwrap_or_default();
processor_tx.send(Message::Exit).unwrap_or_default();
break;
}
};
});
self.tpool.join();
Ok(())
}
}

146
src/server/processor.rs Normal file
View File

@@ -0,0 +1,146 @@
use anyhow::Result;
use std::sync::mpsc::{Receiver, TryRecvError};
use std::sync::Arc;
use super::file::FileMeta;
use super::store::{Index, Store};
use super::{FileQueue, Message};
pub struct Processor {
files: FileQueue,
duplicates: Arc<Store>,
msg_rx: Receiver<Message>,
}
impl Processor {
pub fn new(files: FileQueue, duplicates: Arc<Store>, msg_rx: Receiver<Message>) -> Self {
Self {
files,
duplicates,
msg_rx,
}
}
pub fn process(&self) -> Result<()> {
loop {
match self.msg_rx.try_recv() {
Ok(Message::Exit) => break,
Err(TryRecvError::Empty) => {},
_ => {}
}
let next_file = {
let mut cfiles = self.files.lock().unwrap();
cfiles.pop()
};
match next_file {
None => continue,
Some(file_res) => match FileMeta::new(file_res) {
Err(_) => continue,
Ok(fm) => {
let fm_arc = Arc::new(fm);
self.duplicates
.add(Index::Size(fm_arc.size), fm_arc.clone());
}
},
}
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use anyhow::Result;
use std::sync::mpsc;
use std::sync::{Arc, Mutex};
use std::thread;
use tempfile::TempDir;
use std::fs::File;
use std::io::Write;
#[test]
fn processor_differentiates_files_with_different_sizes() -> Result<()> {
let file_queue = Arc::new(Mutex::new(vec![]));
let (tx, rx) = mpsc::channel::<Message>();
let root = TempDir::new()?;
let store = Arc::new(Store::new());
let fq_c = file_queue.clone();
let store_c = store.clone();
let proc_thread = thread::spawn(move || {
let processor = Processor::new(fq_c, store_c, rx);
processor.process().expect("processor failed");
});
let files =
[("hello.txt", "Lorem ipsum dolor sit amet, consectetur adipiscing elit, sed do eiusmod tempor incididunt ut labore et dolore magna aliqua. Ut enim ad minim veniam, quis nostrud exercitation ullamco laboris nisi ut aliquip ex ea commodo consequat"),
("hello_dup.txt", "Duis aute irure dolor in reprehenderit in voluptate velit esse cillum dolore eu fugiat nulla pariatur. Excepteur sint occaecat cupidatat non proident, sunt in culpa qui officia deserunt mollit anim id est laborum")];
for (filename, content) in files.into_iter() {
let fpath = root.path().join(filename);
let mut tf = File::create(fpath.clone())?;
tf.write_all(content.as_bytes())?;
let mut mfq = file_queue.lock().unwrap();
mfq.push(fpath.into_os_string().into_string().unwrap().into_boxed_str());
}
for _ in 0..10 {
if store.entries().len() < 2 {
thread::sleep(std::time::Duration::from_millis(100));
}
}
tx.send(Message::Exit).expect("unable to send msg to processor");
proc_thread.join().expect("failed to join on thread!");
assert!(store.entries().len() == 2);
Ok(())
}
#[test]
fn processor_groups_files_with_same_size() -> Result<()> {
let file_queue = Arc::new(Mutex::new(vec![]));
let (tx, rx) = mpsc::channel::<Message>();
let root = TempDir::new()?;
let store = Arc::new(Store::new());
let fq_c = file_queue.clone();
let store_c = store.clone();
let proc_thread = thread::spawn(move || {
let processor = Processor::new(fq_c, store_c, rx);
processor.process().expect("processor failed");
});
let files =
[("hello.txt", "Lorem ipsum dolor sit amet, consectetur adipiscing elit, sed do eiusmod tempor incididunt ut labore et dolore magna aliqua. Ut enim ad minim veniam, quis nostrud exercitation ullamco laboris nisi ut aliquip ex ea commodo consequat. Duis aute irure dolor in reprehenderit in voluptate velit esse cillum dolore eu fugiat nulla pariatur. Excepteur sint occaecat cupidatat non proident, sunt in culpa qui officia deserunt mollit anim id est laborum"),
("hello_dup.txt", "Lorem ipsum dolor sit amet, consectetur adipiscing elit, sed do eiusmod tempor incididunt ut labore et dolore magna aliqua. Ut enim ad minim veniam, quis nostrud exercitation ullamco laboris nisi ut aliquip ex ea commodo consequat. Duis aute irure dolor in reprehenderit in voluptate velit esse cillum dolore eu fugiat nulla pariatur. Excepteur sint occaecat cupidatat non proident, sunt in culpa qui officia deserunt mollit anim id est laborum")];
for (filename, content) in files.into_iter() {
let fpath = root.path().join(filename);
let mut tf = File::create(fpath.clone())?;
tf.write_all(content.as_bytes())?;
let mut mfq = file_queue.lock().unwrap();
mfq.push(fpath.into_os_string().into_string().unwrap().into_boxed_str());
}
for _ in 0..10 {
if store.entries().is_empty() {
thread::sleep(std::time::Duration::from_millis(100));
}
}
tx.send(Message::Exit).expect("unable to send msg to processor");
proc_thread.join().expect("failed to join on thread!");
assert!(store.entries().len() == 1);
Ok(())
}
}

128
src/server/scanner.rs Normal file
View File

@@ -0,0 +1,128 @@
use super::{FileQueue, Message};
use anyhow::Result;
use std::fs;
use std::sync::mpsc::{self, Receiver};
use std::sync::Arc;
use std::sync::Mutex;
#[derive(Debug)]
pub struct Scanner {
files: FileQueue,
proc_queue: FileQueue,
msg_rx: Receiver<Message>,
}
impl Scanner {
pub fn new(fq: FileQueue, rx: Receiver<Message>) -> Self {
Self {
files: fq,
proc_queue: Arc::new(Mutex::new(vec![])),
msg_rx: rx,
}
}
pub fn index(&self) -> Result<()> {
loop {
match self.msg_rx.try_recv() {
Ok(Message::AddScanDirectory(path)) => {
let mut mfq = self.proc_queue.lock().unwrap();
mfq.push(path);
}
Err(mpsc::TryRecvError::Empty) | Ok(Message::None) => {}
Ok(Message::Exit) | Err(_) => break,
}
let npath = {
match self.proc_queue.try_lock() {
Ok(mut q) => q.pop(),
Err(_) => None,
}
};
match npath {
None => continue,
Some(path) => {
std::fs::read_dir(path.as_ref())?
.filter_map(Result::ok)
.for_each(|entry: fs::DirEntry| {
let mdata =
fs::metadata(entry.path()).expect("unable to read file metadata.");
let mpath = entry
.path()
.into_os_string()
.into_string()
.expect("invalid path conversion failed.")
.into_boxed_str();
match mdata.is_dir() {
true => {
let mut pq = self
.proc_queue
.lock()
.expect("proc queue lock acq failed.");
pq.push(mpath);
}
false => {
let mut fq = self
.files
.lock()
.expect("file queue lock acq failed.");
fq.push(mpath)
}
}
});
}
}
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use anyhow::Result;
use std::fs::File;
use std::io::Write;
use std::sync::mpsc::channel;
use tempfile::TempDir;
use std::thread;
#[test]
fn scanner_scans_files() -> Result<()> {
let files =
[("hello.txt", "Lorem ipsum dolor sit amet, consectetur adipiscing elit, sed do eiusmod tempor incididunt ut labore et dolore magna aliqua. Ut enim ad minim veniam, quis nostrud exercitation ullamco laboris nisi ut aliquip ex ea commodo consequat. Duis aute irure dolor in reprehenderit in voluptate velit esse cillum dolore eu fugiat nulla pariatur. Excepteur sint occaecat cupidatat non proident, sunt in culpa qui officia deserunt mollit anim id est laborum"),
("hello_dup.txt", "Lorem ipsum dolor sit amet, consectetur adipiscing elit, sed do eiusmod tempor incididunt ut labore et dolore magna aliqua. Ut enim ad minim veniam, quis nostrud exercitation ullamco laboris nisi ut aliquip ex ea commodo consequat. Duis aute irure dolor in reprehenderit in voluptate velit esse cillum dolore eu fugiat nulla pariatur. Excepteur sint occaecat cupidatat non proident, sunt in culpa qui officia deserunt mollit anim id est laborum")];
let file_queue: FileQueue = Arc::new(Mutex::new(vec![]));
let (tx, rx) = channel::<Message>();
let root = TempDir::new()?;
for (filename, content) in files.into_iter() {
let mut tf = File::create(root.path().join(filename))?;
tf.write_all(content.as_bytes())?;
}
let rpath = root.path().to_str().unwrap().to_string().into_boxed_str();
let fq_clone = file_queue.clone();
let scanner_thread = thread::spawn(move || {
let scanner = Scanner::new(fq_clone, rx);
scanner.index().unwrap();
});
tx.send(Message::AddScanDirectory(rpath))?;
tx.send(Message::Exit)?;
scanner_thread.join().unwrap();
let fq_len = {
let v = file_queue.lock().unwrap();
v.len()
};
assert_eq!(files.len(), fq_len);
Ok(())
}
}

35
src/server/store.rs Normal file
View File

@@ -0,0 +1,35 @@
use super::file::FileMeta;
use dashmap::DashMap;
use std::sync::Arc;
#[allow(unused)]
#[derive(Debug, Hash, PartialEq, Eq)]
pub enum Index {
Size(u64),
Partial(Arc<[u8]>),
Full(Box<str>),
}
#[derive(Debug)]
pub struct Store {
internal: Arc<DashMap<Index, Vec<Arc<FileMeta>>>>,
}
impl Store {
pub fn new() -> Self {
Self {
internal: Arc::new(DashMap::new()),
}
}
pub fn entries(&self) -> Arc<DashMap<Index, Vec<Arc<FileMeta>>>> {
self.internal.clone()
}
pub fn add(&self, index: Index, file: Arc<FileMeta>) {
self.internal
.entry(index)
.and_modify(|fg| fg.push(file.clone()))
.or_insert(vec![file]);
}
}

99
src/tui/mod.rs Normal file
View File

@@ -0,0 +1,99 @@
use std::sync::mpsc::Sender;
use std::time::Duration;
use anyhow::{Context, Result};
use ratatui::crossterm::event::{self, Event, KeyCode};
use ratatui::layout::{Constraint, Direction, Layout};
use ratatui::widgets::{Block, Borders, List, ListItem, Paragraph};
use ratatui::{DefaultTerminal, Frame};
use crate::server::{Message, Server};
use std::sync::Arc;
pub struct Tui {
app_tx: Sender<Message>,
server: Arc<Server>,
}
impl Tui {
pub fn new(app_tx: Sender<Message>, server: Arc<Server>) -> Self {
Self { app_tx, server }
}
pub fn start(&mut self) -> Result<()> {
let terminal = ratatui::init();
self.run(terminal).expect("ui loop failed.");
ratatui::restore();
Ok(())
}
fn poll_events() -> Result<Message> {
match event::poll(Duration::from_millis(100)).context("event polling failed.")? {
true => match event::read().context("event read failed.")? {
Event::Key(key) => match key.code {
KeyCode::Char('q') => Ok(Message::Exit),
_ => Ok(Message::None),
},
_ => Ok(Message::None),
},
false => Ok(Message::None),
}
}
fn handle_events(&self) -> Result<Message> {
match Self::poll_events() {
Ok(Message::Exit) => {
self.app_tx
.send(Message::Exit)
.expect("app event send failed.");
Ok(Message::Exit)
}
_ => Ok(Message::None),
}
}
fn draw(&mut self, frame: &mut Frame) {
let listitems = self
.server
.dupstore
.entries()
.iter()
.flat_map(|val| {
let mut group = vec![ListItem::new(format!("{:?}", val.key()))];
val.value().iter().for_each(|f| {
group.push(ListItem::new(format!("|-{}", f.path)));
});
group
})
.collect::<Vec<ListItem>>();
let list = List::new(listitems.clone());
let layout_chunks = Layout::default()
.direction(Direction::Horizontal)
.constraints([Constraint::Percentage(50), Constraint::Percentage(50)])
.split(frame.area());
frame.render_widget(
list.block(Block::new().borders(Borders::ALL)),
layout_chunks[0],
);
frame.render_widget(
Paragraph::new(String::new()).block(Block::new().borders(Borders::ALL)),
layout_chunks[1],
);
}
fn run(&mut self, mut terminal: DefaultTerminal) -> Result<()> {
loop {
terminal.draw(|f| self.draw(f))?;
match self.handle_events() {
Ok(Message::Exit) => break,
_ => {}
}
}
Ok(())
}
}