summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorjakka <jakkadoujin@gmail.com>2025-06-12 21:48:22 +0300
committerjakka <jakkadoujin@gmail.com>2025-06-12 21:48:22 +0300
commit743b87fc7aea1aa90cfacee17cdac8ca818c5661 (patch)
tree0b1edf14b68a242a964b6a6704fb55c53c66dc4e
parent382e526f76ada0ca58f560a1217cb9a597ba3945 (diff)
continued writing db logic. added files indexing
-rw-r--r--Cargo.toml2
-rw-r--r--src/db.rs139
-rw-r--r--src/files.rs135
-rw-r--r--src/flac.rs7
-rw-r--r--src/main.rs1
5 files changed, 238 insertions, 46 deletions
diff --git a/Cargo.toml b/Cargo.toml
index 2145830..cc8e538 100644
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -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"] }
diff --git a/src/db.rs b/src/db.rs
index 9346575..cb63575 100644
--- a/src/db.rs
+++ b/src/db.rs
@@ -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;