diff options
| author | jakka <jakkadoujin@gmail.com> | 2025-06-12 21:48:22 +0300 |
|---|---|---|
| committer | jakka <jakkadoujin@gmail.com> | 2025-06-12 21:48:22 +0300 |
| commit | 743b87fc7aea1aa90cfacee17cdac8ca818c5661 (patch) | |
| tree | 0b1edf14b68a242a964b6a6704fb55c53c66dc4e | |
| parent | 382e526f76ada0ca58f560a1217cb9a597ba3945 (diff) | |
continued writing db logic. added files indexing
| -rw-r--r-- | Cargo.toml | 2 | ||||
| -rw-r--r-- | src/db.rs | 139 | ||||
| -rw-r--r-- | src/files.rs | 135 | ||||
| -rw-r--r-- | src/flac.rs | 7 | ||||
| -rw-r--r-- | src/main.rs | 1 |
5 files changed, 238 insertions, 46 deletions
@@ -14,5 +14,5 @@ libsql = { version = "0.9.10", default-features = false, features = ["core", "sy md-5 = "0.10.6" metaflac = "0.2.8" symphonia = { version = "0.5.4", path = "../Symphonia/symphonia", default-features = false, features = ["flac"] } -tokio = { version = "1.45.1", features = ["macros", "rt-multi-thread"] } +tokio = { version = "1.45.1", features = ["macros", "rt", "rt-multi-thread"] } #symphonia = { git = "https://github.com/pdeljanov/Symphonia.git", branch = "dev-0.6", default-features = false, features = ["flac"] } @@ -1,4 +1,4 @@ -use anyhow::{Ok, Result, anyhow}; +use anyhow::{Result, anyhow}; use directories::BaseDirs; use futures_util::StreamExt; use libsql::{Builder, Connection, params}; @@ -6,22 +6,19 @@ use std::{ ffi::OsStr, fmt::Display, path::{Path, absolute}, + time::{Duration, UNIX_EPOCH}, }; -use crate::flac::{get_vendor, CURRENT_VENDOR}; +use crate::flac::{CURRENT_VENDOR, get_vendor}; -pub async fn open_db() -> Result<Connection> { - let conn = if let Some(base_dir) = BaseDirs::new() { - let db_name = Path::new(base_dir.data_dir()).join("reencoder.db"); - Builder::new_local(db_name).build().await?.connect()? - } else { - return Err(anyhow!("Failed to locate data directory")); - }; - - conn.execute("CREATE TABLE IF NOT EXISTS flacs (path TEXT PRIMARY KEY, toencode BOOLEAN NOT NULL)", ()).await?; - - Ok(conn) -} +const TABLE_CREATE: &str = "CREATE TABLE IF NOT EXISTS flacs (path TEXT PRIMARY KEY, toencode BOOLEAN NOT NULL, modtime INTEGER)"; +const ADD_NEW_ITEM: &str = "INSERT INTO flacs (path, toencode, modtime) VALUES (?1, ?2, ?3)"; +const REPLACE_ITEM: &str = "REPLACE INTO flacs (path, toencode, modtime) VALUES (?1, ?2, ?3)"; +const TOENCODE_QUERY: &str = "SELECT path FROM flacs WHERE toencode"; +const CHECK_FILE: &str = "SELECT exists(SELECT 1 FROM flacs WHERE path = ?1)"; +const FETCH_MODTIME: &str = "SELECT modtime FROM flacs WHERE path = ?1"; +const FETCH_FILES: &str = "SELECT path FROM flacs"; +const REMOVE_FILE: &str = "DELETE FROM flac WHERE path = ?1"; #[derive(Debug)] pub enum Errors { @@ -40,19 +37,25 @@ pub trait Reencoder { async fn insert_file(&self, filename: &impl AsRef<OsStr>) -> Result<()>; async fn update_file(&self, filename: &impl AsRef<OsStr>) -> Result<()>; async fn get_files_toencode(&self) -> Result<Vec<String>>; - async fn get_files_indexed(&self) -> Result<Vec<String>>; + async fn check_file(&self, filename: &impl AsRef<OsStr>) -> Result<bool>; + async fn get_modtime(&self, filename: &impl AsRef<OsStr>) -> Result<u64>; + async fn clean_files(&self) -> Result<()>; } impl Reencoder for Connection { async fn insert_file(&self, filename: &impl AsRef<OsStr>) -> Result<()> { - let file = Path::new(filename); - let toencode = !matches!(get_vendor(file).as_str(), CURRENT_VENDOR); + let abs_filename = absolute(Path::new(filename))?; + let toencode = !matches!(get_vendor(&abs_filename)?.as_str(), CURRENT_VENDOR); - let abs_filename = absolute(file)?; + let modtime = abs_filename + .metadata()? + .modified()? + .duration_since(UNIX_EPOCH)? + .as_secs(); self.execute( - "INSERT INTO flacs (path, toencode) VALUES (?1, ?2)", - params![abs_filename.to_str().unwrap(), toencode], + ADD_NEW_ITEM, + params![abs_filename.to_str().unwrap(), toencode, modtime], ) .await?; @@ -62,9 +65,15 @@ impl Reencoder for Connection { async fn update_file(&self, filename: &impl AsRef<OsStr>) -> Result<()> { let abs_filename = absolute(Path::new(filename))?; + let modtime = abs_filename + .metadata()? + .modified()? + .duration_since(UNIX_EPOCH)? + .as_secs(); + self.execute( - "REPLACE INTO flacs (path, toencode) VALUES (?1, ?2)", - params![abs_filename.to_str().unwrap(), false], + REPLACE_ITEM, + params![abs_filename.to_str().unwrap(), false, modtime], ) .await?; @@ -72,9 +81,7 @@ impl Reencoder for Connection { } async fn get_files_toencode(&self) -> Result<Vec<String>> { - let rows = self - .query("SELECT path FROM flacs WHERE toencode", ()) - .await?; + let rows = self.query(TOENCODE_QUERY, ()).await?; if rows.column_count() == 0 { return Err(anyhow!(Errors::EmptyQuery)); }; @@ -88,24 +95,73 @@ impl Reencoder for Connection { Ok(filenames) } - async fn get_files_indexed(&self) -> Result<Vec<String>> { - let rows = self - .query("SELECT path FROM flacs", ()) - .await?; - if rows.column_count() == 0 { - return Err(anyhow!(Errors::EmptyQuery)); - }; + async fn check_file(&self, filename: &impl AsRef<OsStr>) -> Result<bool> { + let abs_filename = absolute(Path::new(filename))?; - let filenames = rows - .into_stream() - .map(|row| row.unwrap().get_str(0).unwrap().to_string()) - .collect::<Vec<String>>() - .await; + if let Some(row) = self + .query(CHECK_FILE, params!(abs_filename.to_str().unwrap())) + .await? + .next() + .await? + { + Ok(matches!(row.get_value(0)?, libsql::Value::Integer(1))) + } else { + Err(anyhow!("database error")) + } + } - Ok(filenames) + async fn get_modtime(&self, filename: &impl AsRef<OsStr>) -> Result<u64> { + let abs_filename = absolute(Path::new(filename))?; + + if let Some(row) = self + .query(FETCH_MODTIME, params!(abs_filename.to_str().unwrap())) + .await? + .next() + .await? + { + if let Some(sec) = row.get_value(0)?.as_integer() { + Ok(Duration::from_secs(*sec as u64).as_secs()) + } else { + Ok(Duration::from_secs(0).as_secs()) + } + } else { + Err(anyhow!("database error")) + } + } + + async fn clean_files(&self) -> Result<()> { + let mut tasks = tokio::task::JoinSet::new(); + while let Ok(Some(row)) = self.query(FETCH_FILES, ()).await?.next().await { + let path = absolute(Path::new(row.get_str(0)?))?; + let conn = self.clone(); + tasks.spawn(async move { + if !path.exists() { + let _ = conn + .execute(REMOVE_FILE, params!(path.to_str().unwrap())) + .await; + } + }); + } + + tasks.join_all().await; + + Ok(()) } } +pub async fn open_db() -> Result<Connection> { + let conn = if let Some(base_dir) = BaseDirs::new() { + let db_name = Path::new(base_dir.data_dir()).join("reencoder.db"); + Builder::new_local(db_name).build().await?.connect()? + } else { + return Err(anyhow!("Failed to locate data directory")); + }; + + conn.execute(TABLE_CREATE, ()).await?; + + Ok(conn) +} + #[cfg(test)] mod tests { use super::*; @@ -117,7 +173,7 @@ mod tests { .unwrap() .connect() .unwrap(); - conn.execute("CREATE TABLE IF NOT EXISTS flacs (path TEXT PRIMARY KEY, toencode BOOLEAN NOT NULL)", ()).await.unwrap(); + conn.execute(TABLE_CREATE, ()).await.unwrap(); conn } @@ -145,10 +201,11 @@ mod tests { let _ = conn .execute( - "REPLACE INTO flacs (path,, toencode) VALUES (?1, ?2)", + REPLACE_ITEM, params![ absolute(Path::new("16bit.flac")).unwrap().to_str(), - true + true, + "" ], ) .await; diff --git a/src/files.rs b/src/files.rs new file mode 100644 index 0000000..df596f8 --- /dev/null +++ b/src/files.rs @@ -0,0 +1,135 @@ +use anyhow::{Result, anyhow}; +use libsql::Connection; +use std::{ + fmt::Display, + path::{Path, PathBuf, absolute}, + time::UNIX_EPOCH, +}; +use tokio::{fs::read_dir, task::JoinSet}; + +use crate::db::Reencoder; + +#[derive(Debug)] +struct FileError { + file: PathBuf, + error: anyhow::Error, +} + +impl Display for FileError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!( + f, + "error: {}\t on file {}", + self.error, + self.file.to_string_lossy() + ) + } +} + +async fn handle_file(file: PathBuf, conn: Connection) -> Result<()> { + match conn.check_file(&file).await { + Ok(true) => { + let modtime = file + .metadata()? + .modified()? + .duration_since(UNIX_EPOCH)? + .as_secs(); + let db_time = conn.get_modtime(&file).await?; + if modtime != db_time { + if let Err(error) = conn.update_file(&file).await { + return Err(anyhow!(FileError { file, error })); + }; + } + return Ok(()); + } + Err(error) => return Err(anyhow!(FileError { file, error })), + _ => {} + } + + if let Err(error) = conn.insert_file(&file).await { + return Err(anyhow!(FileError { file, error })); + } + + Ok(()) +} + +pub async fn index_files_recursively(path: &Path, conn: &Connection) -> Result<()> { + if !path.is_dir() { + return Err(anyhow!("Invalid root directory")); + } + let abspath = absolute(path)?; + let mut tasks = JoinSet::new(); + let mut counter: i64 = 0; + + let mut dirs = vec![abspath]; + + while let Some(dir) = dirs.pop() { + let mut read_dir = read_dir(dir).await?; + + while let Some(entry) = read_dir.next_entry().await? { + let path = entry.path(); + if path.is_dir() { + dirs.push(path); + } else if path.is_file() { + if let Some(ext) = path.extension() { + if ext == "flac" { + let newconn = conn.clone(); + counter += 1; + tasks.spawn(async move { handle_file(path, newconn).await }); + } + } + } + print!("\rFiles found:\t{counter}") + } + } + + while let Some(task) = tasks.join_next().await { + match task { + Ok(Err(error)) => eprintln!("{error}"), + Err(error) => eprintln!("Error encountered:\t{}", error), + _ => {} + } + } + + Ok(()) +} + +#[cfg(test)] +mod tests { + use libsql::Builder; + + use super::*; + + async fn dummy_db(name: impl AsRef<Path>) -> Connection { + let conn = Builder::new_local(name) + .build() + .await + .unwrap() + .connect() + .unwrap(); + conn.execute("CREATE TABLE IF NOT EXISTS flacs (path TEXT PRIMARY KEY, toencode BOOLEAN NOT NULL, modtime INTEGER)", ()).await.unwrap(); + conn + } + + #[tokio::test] + async fn test_lots_of_files() { + let conn = dummy_db("temp3.db").await; + index_files_recursively(Path::new("/mnt/Music"), &conn) + .await + .unwrap(); + println!( + "\n{}", + conn.query("SELECT COUNT(DISTINCT path) FROM flacs", ()) + .await + .unwrap() + .next() + .await + .unwrap() + .unwrap() + .get_value(0) + .unwrap() + .as_integer() + .unwrap() + ); + } +} diff --git a/src/flac.rs b/src/flac.rs index 04504f8..16627e2 100644 --- a/src/flac.rs +++ b/src/flac.rs @@ -12,7 +12,6 @@ use symphonia::core::{ meta::MetadataOptions, }; -#[allow(dead_code)] pub const CURRENT_VENDOR: &str = "reference libFLAC 1.5.0 20250211"; struct StreamConfig { @@ -280,9 +279,9 @@ pub fn encode_file(filename: impl AsRef<OsStr>) -> Result<()> { Ok(()) } -pub fn get_vendor(file: &Path) -> String { - let tag = Tag::read_from_path(file).unwrap(); - tag.vorbis_comments().unwrap().vendor_string.clone() +pub fn get_vendor(file: &Path) -> Result<String> { + let tag = Tag::read_from_path(file)?; + Ok(tag.vorbis_comments().unwrap().vendor_string.clone()) } #[cfg(test)] diff --git a/src/main.rs b/src/main.rs index 05e3bc2..00a00c7 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,4 +1,5 @@ mod db; +mod files; mod flac; use anyhow::Result; |
