summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--src/files.rs74
1 files changed, 34 insertions, 40 deletions
diff --git a/src/files.rs b/src/files.rs
index 83d445c..3af2010 100644
--- a/src/files.rs
+++ b/src/files.rs
@@ -141,54 +141,48 @@ pub fn reencode_files(conn: Connection, handler: Arc<AtomicBool>, threads: usize
let mut files = conn.get_toencode_files()?.into_iter();
+ let lock = Arc::new(Mutex::new(conn));
+
let thread_counter = Arc::new(AtomicUsize::new(0));
- let lock = Arc::new(Mutex::new(conn));
+ thread::scope(|s| {
+ while handler.load(Ordering::SeqCst) {
+ if thread_counter.load(Ordering::Relaxed) >= threads {
+ sleep(Duration::from_millis(100));
+ #[cfg(not(test))]
+ bar.tick();
+ continue;
+ }
- while handler.load(Ordering::SeqCst) {
- if thread_counter.load(Ordering::Relaxed) >= threads {
- sleep(Duration::from_millis(100));
- continue;
- }
+ let file = match files.next() {
+ Some(file) => file,
+ None => break,
+ };
- let file = match files.next() {
- Some(file) => file,
- None => break,
- };
+ thread_counter.fetch_add(1, Ordering::Relaxed);
- let newhandler = handler.clone();
- let newlock = lock.clone();
- #[cfg(not(test))]
- let newbar = bar.clone();
- let newcounter = thread_counter.clone();
- thread_counter.fetch_add(1, Ordering::Relaxed);
+ let handler = handler.clone();
+ let lock = lock.clone();
+ #[cfg(not(test))]
+ let bar = bar.clone();
+ let thread_counter = thread_counter.clone();
- let _ = thread::spawn(move || {
- match handle_encode(&file, newhandler) {
- Err(error) => eprintln!("{}", FileError::new(&file, error)),
- Ok(false) => {
- let conn = match newlock.lock() {
- Ok(conn) => conn,
- Err(_) => {
- eprintln!("Lock poisoned on file:\t{}", file.to_string_lossy());
- return;
+ s.spawn(move || {
+ match handle_encode(&file, handler) {
+ Err(error) => eprintln!("{}", FileError::new(&file, error)),
+ Ok(false) => {
+ if let Err(error) = lock.lock().unwrap().update_file(&file) {
+ eprintln!("{}", FileError::new(file, error));
}
- };
- if let Err(error) = conn.update_file(&file) {
- eprintln!("{}", FileError::new(file, error));
+ #[cfg(not(test))]
+ bar.inc(1)
}
- #[cfg(not(test))]
- newbar.inc(1)
- }
- Ok(true) => {}
- };
- newcounter.fetch_sub(1, Ordering::Relaxed);
- });
- }
-
- while thread_counter.load(Ordering::Relaxed) != 0 {
- sleep(Duration::from_millis(100));
- }
+ Ok(true) => {}
+ };
+ thread_counter.fetch_sub(1, Ordering::Relaxed);
+ });
+ }
+ });
#[cfg(not(test))]
{