From 5ce0a2774aee953352c8984b016ef0422cd1108b Mon Sep 17 00:00:00 2001 From: sreedevk Date: Mon, 7 Jul 2025 00:00:25 +0000 Subject: [PATCH] added server threadpool and message passing --- Cargo.lock | 4 +- src/cli/mod.rs | 0 src/main.rs | 4 ++ src/pipeline/mod.rs | 93 +++++++++++++++++++++++++++------------ src/pipeline/processor.rs | 37 ++++++++++++++++ src/pipeline/scanner.rs | 71 +++++++++++++++--------------- src/pipeline/store.rs | 32 ++++++++++++++ src/tui/mod.rs | 0 8 files changed, 175 insertions(+), 66 deletions(-) create mode 100644 src/cli/mod.rs create mode 100644 src/pipeline/processor.rs create mode 100644 src/pipeline/store.rs create mode 100644 src/tui/mod.rs diff --git a/Cargo.lock b/Cargo.lock index 5aa2658..c3c50b5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -538,9 +538,9 @@ checksum = "bbd2bcb4c963f2ddae06a2efc7e9f3591312473c50c6685e1f298068316e66fe" [[package]] name = "libc" -version = "0.2.155" +version = "0.2.174" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "97b3888a4aecf77e811145cadf6eef5901f4782c53886191b2f693f24761847c" +checksum = "1171693293099992e19cddea4e8b849964e9846f4acee11b3948bcc337be8776" [[package]] name = "libredox" diff --git a/src/cli/mod.rs b/src/cli/mod.rs new file mode 100644 index 0000000..e69de29 diff --git a/src/main.rs b/src/main.rs index 62dc181..04193b7 100644 --- a/src/main.rs +++ b/src/main.rs @@ -4,7 +4,11 @@ mod interactive; mod params; mod processor; mod scanner; + +/* version 2.0 modules*/ mod pipeline; +mod cli; + use anyhow::Result; use clap::Parser; diff --git a/src/pipeline/mod.rs b/src/pipeline/mod.rs index e86bdd0..bb78f4f 100644 --- a/src/pipeline/mod.rs +++ b/src/pipeline/mod.rs @@ -1,41 +1,76 @@ mod file; +mod processor; mod scanner; +mod store; + +use anyhow::{anyhow, Result}; +use std::collections::HashMap; +use std::sync::mpsc; +use std::sync::mpsc::channel; use std::sync::Arc; use std::sync::Mutex; +use threadpool::ThreadPool; use self::file::FileMeta; +use self::processor::Processor; +use self::scanner::Scanner; +use self::store::Store; -const CONCURRENCY: usize = 4; +pub type FileQueue = Arc>>>; -pub struct Server { } +pub enum Message { + AddScanDirectory(Box), + Exit, +} + +pub struct Server { + fq: FileQueue, + dupstore: Arc, + tpool: ThreadPool, +} impl Server { - // pub fn new() -> Result { - // Ok(Self { - // index_queue: Arc::new(Mutex::new(vec![])), - // index_tpool: ThreadPool::new(CONCURRENCY), - // process_queue: Arc::new(Mutex::new(vec![])), - // }) - // } + pub fn new() -> Result { + Ok(Self { + fq: Arc::new(Mutex::new(vec![])), + dupstore: Arc::new(Store::new()), + tpool: ThreadPool::new(4), + }) + } - // fn index(&self) -> Result<()> { - // let mut iq = self.index_queue.lock().unwrap(); - // let mut pq = self.process_queue.lock().unwrap(); - // - // match iq.pop() { - // None => Ok(()), - // Some(QueueElem::File(path)) => { - // pq.push(FileMeta::new(path)?); - // - // Ok(()) - // } - // Some(QueueElem::Directory(path)) => { - // self.index_tpool.execute(move || { - // Scanner::new(&path); - // }); - // - // Ok(()) - // } - // } - // } + pub fn start(&self, rx: mpsc::Receiver) -> Result<()> { + let processor_fq = self.fq.clone(); + let processor_store = self.dupstore.clone(); + let (processor_tx, processor_rx) = channel::(); + self.tpool.execute(move || { + Processor::new(processor_fq, processor_store, processor_rx) + .process() + .ok(); + }); + + let scanner_fq = self.fq.clone(); + let (scanner_tx, scanner_rx) = channel::(); + self.tpool.execute(move || { + Scanner::new(scanner_fq, scanner_rx).index().ok(); + }); + + self.tpool.execute(move || loop { + match rx.recv() { + Ok(Message::AddScanDirectory(path)) => { + scanner_tx + .send(Message::AddScanDirectory(path)) + .unwrap_or_default(); + } + 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(()) + } } diff --git a/src/pipeline/processor.rs b/src/pipeline/processor.rs new file mode 100644 index 0000000..fb12e65 --- /dev/null +++ b/src/pipeline/processor.rs @@ -0,0 +1,37 @@ +use anyhow::Result; +use std::collections::HashMap; +use std::sync::mpsc::Receiver; +use std::sync::{Arc, Mutex}; + +use super::file::FileMeta; +use super::Store; +use super::{FileQueue, Message}; + +const BATCH_SIZE: usize = 10; + +pub struct Processor { + files: FileQueue, + duplicates: Arc, + msg_rx: Receiver, +} + +impl Processor { + pub fn new(files: FileQueue, duplicates: Arc, msg_rx: Receiver) -> Self { + Self { + files, + duplicates, + msg_rx, + } + } + + pub fn process(&self) -> Result<()> { + loop { + match self.msg_rx.try_recv() { + Ok(Message::Exit) | Err(_) => break, + _ => {} + } + } + + Ok(()) + } +} diff --git a/src/pipeline/scanner.rs b/src/pipeline/scanner.rs index 79d9226..64c6001 100644 --- a/src/pipeline/scanner.rs +++ b/src/pipeline/scanner.rs @@ -1,57 +1,59 @@ +use super::{FileQueue, Message}; use anyhow::{anyhow, Result}; use std::fs; use std::sync::atomic::{AtomicU32, Ordering::Relaxed}; +use std::sync::mpsc::{self, Receiver}; use std::sync::Arc; use std::sync::Mutex; use threadpool::ThreadPool; pub struct Scanner { - files: Arc>>>, + files: FileQueue, threadpool: Arc, - proc_count: Arc, - proc_queue: Arc>>>, + proc_queue: FileQueue, + msg_rx: Receiver, } impl Scanner { - pub fn new(path: Box) -> Result { - Ok(Self { - files: Arc::new(Mutex::new(vec![])), + pub fn new(fq: FileQueue, rx: Receiver) -> Self { + Self { + files: fq, threadpool: Arc::new(ThreadPool::new(8)), - proc_count: Arc::new(AtomicU32::new(1)), - proc_queue: Arc::new(Mutex::new(vec![path])), - }) + proc_queue: Arc::new(Mutex::new(vec![])), + msg_rx: rx, + } } pub fn index(&self) -> Result<()> { loop { - let pc = self.proc_count.load(Relaxed); - match pc { - 0 => break, - _ => { - let npath = { - let mut q = self.proc_queue.lock().unwrap(); - q.pop() - }; + 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::Exit) | Err(_) => break, + } - match npath { - None => continue, - Some(path) => { - self.proc_count.fetch_add(1, Relaxed); - let copy_of_queue = self.proc_queue.clone(); - let copy_of_files = self.files.clone(); - let copy_of_proc_count = self.proc_count.clone(); + let npath = { + let mut q = self.proc_queue.lock().unwrap(); + q.pop() + }; - self.threadpool.execute(move || { - if let Ok((files, dirs)) = Self::scan(path) { - let mut q = copy_of_queue.lock().unwrap(); - let mut f = copy_of_files.lock().unwrap(); - q.extend(files); - f.extend(dirs); - copy_of_proc_count.fetch_sub(1, Relaxed); - } - }); + match npath { + None => continue, + Some(path) => { + let copy_of_queue = self.proc_queue.clone(); + let copy_of_files = self.files.clone(); + + self.threadpool.execute(move || { + if let Ok((files, dirs)) = Self::scan(path) { + let mut q = copy_of_queue.lock().unwrap(); + let mut f = copy_of_files.lock().unwrap(); + q.extend(files); + f.extend(dirs); } - } + }); } } } @@ -74,7 +76,6 @@ impl Scanner { .unwrap() .into_boxed_str(); - // TODO: Support Symlinks if entrymeta.is_file() { files.push(path); } else if entrymeta.is_dir() { diff --git a/src/pipeline/store.rs b/src/pipeline/store.rs new file mode 100644 index 0000000..fd78c9d --- /dev/null +++ b/src/pipeline/store.rs @@ -0,0 +1,32 @@ +use super::file::FileMeta; +use anyhow::Result; +use std::collections::HashMap; +use std::sync::{Arc, Mutex}; + +#[derive(Debug, Hash, PartialEq, Eq)] +pub enum Index { + Size(Box), + Partial(Box), + Full(Box), +} + +pub struct Store { + internal: Arc>>>>, +} + +impl Store { + pub fn new() -> Self { + Self { + internal: Arc::new(Mutex::new(HashMap::new())), + } + } + + pub fn add(&self, index: Index, file: Arc) -> Result<()> { + let mut imut = self.internal.lock().unwrap(); + imut.entry(index) + .and_modify(|fg| fg.push(file.clone())) + .or_insert(vec![file]); + + Ok(()) + } +} diff --git a/src/tui/mod.rs b/src/tui/mod.rs new file mode 100644 index 0000000..e69de29