diff options
| author | jakka <jakkadoujin@gmail.com> | 2025-06-19 22:28:40 +0300 |
|---|---|---|
| committer | jakka <jakkadoujin@gmail.com> | 2025-06-19 22:28:40 +0300 |
| commit | 631a904769fb7859adfb7e0ca9fb1bb20f9c5c77 (patch) | |
| tree | 7be6ba6892ec363d437b129df080ac691985b076 /src | |
| parent | 1dde9216b9fc231720815963a5417066f3160248 (diff) | |
checking files before reencodingv0.1.2
Diffstat (limited to 'src')
| -rw-r--r-- | src/files.rs | 78 | ||||
| -rw-r--r-- | src/flac.rs | 7 |
2 files changed, 58 insertions, 27 deletions
diff --git a/src/files.rs b/src/files.rs index bbf4acb..bbcc0c4 100644 --- a/src/files.rs +++ b/src/files.rs @@ -1,6 +1,6 @@ use anyhow::{Result, anyhow}; use futures_util::StreamExt; -#[allow(unused_imports)] +#[cfg(not(test))] use indicatif::{ProgressBar, ProgressStyle}; use pin_utils::pin_mut; use std::{ @@ -15,7 +15,7 @@ use walkdir::WalkDir; use crate::{db::Database, flac::encode_file}; -#[allow(dead_code)] +#[cfg(not(test))] const BAR_TEMPLATE: &str = "{msg} [{wide_bar:.green/cyan}] Elapsed: {elapsed} {pos:>7}/{len:7}"; #[derive(Debug)] @@ -84,7 +84,7 @@ pub async fn index_files_recursively( } let abspath = path.as_ref().canonicalize()?; - let mut tasks = JoinSet::new(); + let mut tasks: JoinSet<Result<(), anyhow::Error>> = JoinSet::new(); #[cfg(not(test))] let bar = ProgressBar::new(0) @@ -99,7 +99,14 @@ pub async fn index_files_recursively( if let Some(ext) = path.extension() { if ext == "flac" { let newconn = conn.clone(); - tasks.spawn(async move { handle_file(path, newconn).await }); + #[cfg(not(test))] + let newbar = bar.clone(); + tasks.spawn(async move { + handle_file(path, newconn).await?; + #[cfg(not(test))] + newbar.inc(1); + Ok(()) + }); #[cfg(not(test))] bar.inc_length(1); } @@ -110,9 +117,18 @@ pub async fn index_files_recursively( while let Some(task) = tokio::select! { _ = canceltoken.cancelled() => { + tasks.abort_all(); + + while let Some(task) = tasks.join_next().await { + match task { + Ok(Err(error)) => eprintln!("{error}"), + Err(error) => if !error.is_cancelled() {eprintln!("Error encountered:\t{}", error)}, + _ => {} + } + } + #[cfg(not(test))] bar.abandon_with_message("Indexing aborted"); - tasks.shutdown().await; return Ok(()) }, task = tasks.join_next() => task @@ -120,10 +136,7 @@ pub async fn index_files_recursively( match task { Ok(Err(error)) => eprintln!("{error}"), Err(error) => eprintln!("Error encountered:\t{}", error), - _ => { - #[cfg(not(test))] - bar.inc(1); - } + _ => {} } } @@ -146,27 +159,43 @@ pub async fn reencode_files(conn: &Database, canceltoken: CancellationToken) -> while let Some(Ok(row)) = stream.next().await { if let Some(file) = row.get_value(0)?.as_text() { let filename = Path::new(file).canonicalize()?; - let newconn = conn.clone(); - tasks.spawn(async move { - let file = filename.clone(); - if let Err(error) = tokio::task::spawn_blocking(move || encode_file(file)).await? { - return Err(anyhow!(FileError::new(&filename, error))); - }; - - if let Err(error) = newconn.update_file(&filename).await { - return Err(anyhow!(FileError::new(&filename, error))); - }; + if filename.exists() { + let newconn = conn.clone(); + #[cfg(not(test))] + let newbar = bar.clone(); + tasks.spawn(async move { + let file = filename.clone(); + if let Err(error) = + tokio::task::spawn_blocking(move || encode_file(file)).await? + { + return Err(anyhow!(FileError::new(&filename, error))); + }; - Ok(()) - }); + if let Err(error) = newconn.update_file(&filename).await { + return Err(anyhow!(FileError::new(&filename, error))); + }; + #[cfg(not(test))] + newbar.inc(1); + Ok(()) + }); + } } } while let Some(task) = tokio::select! { _ = canceltoken.cancelled() => { + tasks.abort_all(); + + while let Some(task) = tasks.join_next().await { + match task { + Ok(Err(error)) => eprintln!("{error}"), + Err(error) => if !error.is_cancelled() {eprintln!("Error encountered:\t{}", error)}, + _ => {} + } + } + #[cfg(not(test))] bar.abandon_with_message("Reencoding aborted"); - tasks.shutdown().await; return Ok(()) }, task = tasks.join_next() => task @@ -174,10 +203,7 @@ pub async fn reencode_files(conn: &Database, canceltoken: CancellationToken) -> match task { Ok(Err(error)) => eprintln!("Error encountered:\t{error}"), Err(error) => eprintln!("Error encountered:\t{error}"), - _ => { - #[cfg(not(test))] - bar.inc(1); - } + _ => {} } } diff --git a/src/flac.rs b/src/flac.rs index 2870500..4283aad 100644 --- a/src/flac.rs +++ b/src/flac.rs @@ -283,8 +283,13 @@ fn encode_cycle_32( pub fn encode_file(filename: impl AsRef<Path>) -> Result<()> { let filencoder = FileEncoder::new(filename)?; + let temp_name = filencoder.temp_name(); - let mut outf = File::create(filencoder.temp_name())?; + if temp_name.exists() { + std::fs::remove_file(&temp_name)?; + } + + let mut outf = File::create(temp_name)?; let mut outw = WriteWrapper(&mut outf); let enc = FlacEncoder::new() .unwrap() |
