diff --git a/Cargo.lock b/Cargo.lock index 82135fc0..cb6f41ea 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8,7 +8,7 @@ version = "0.5.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5f7b0a21988c1bf877cf4759ef5ddaac04c1c9fe808c9142ecb78ba97d97a28a" dependencies = [ - "bitflags", + "bitflags 2.8.0", "bytes", "futures-core", "futures-sink", @@ -31,7 +31,7 @@ dependencies = [ "actix-utils", "ahash 0.8.11", "base64", - "bitflags", + "bitflags 2.8.0", "brotli", "bytes", "bytestring", @@ -508,7 +508,7 @@ version = "54.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0f40f6be8f78af1ab610db7d9b236e21d587b7168e368a36275d2e5670096735" dependencies = [ - "bitflags", + "bitflags 2.8.0", ] [[package]] @@ -665,6 +665,12 @@ version = "1.6.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8c3c1a368f70d6cf7302d78f8f7093da241fb8e8807c05cc9e51a125895a6d5b" +[[package]] +name = "bitflags" +version = "1.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a" + [[package]] name = "bitflags" version = "2.8.0" @@ -1114,10 +1120,10 @@ version = "0.28.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "829d955a0bb380ef178a640b91779e3987da38c9aea133b20614cfed8cdea9c6" dependencies = [ - "bitflags", + "bitflags 2.8.0", "crossterm_winapi", "parking_lot", - "rustix", + "rustix 0.38.44", "winapi", ] @@ -1267,14 +1273,20 @@ dependencies = [ "chrono", "ctr", "dotenv", + "futures", "hex", "jsonwebtoken", + "lofty", + "md5", "owo-colors", "redis", "reqwest", "serde", "serde_json", + "sha256", "sqlx", + "symphonia", + "tempfile", "tokio", "tokio-stream", ] @@ -1401,6 +1413,12 @@ dependencies = [ "pin-project-lite", ] +[[package]] +name = "extended" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "af9673d8203fcb076b19dfd17e38b3d4ae9f44959416ea532ce72415a6020365" + [[package]] name = "fallible-iterator" version = "0.3.0" @@ -1649,15 +1667,21 @@ dependencies = [ "chrono", "ctr", "dotenv", + "futures", "hex", "jsonwebtoken", + "lofty", + "md5", "owo-colors", "redis", "reqwest", "serde", "serde_json", "serde_urlencoded", + "sha256", "sqlx", + "symphonia", + "tempfile", "tokio", "tokio-stream", ] @@ -2254,7 +2278,7 @@ version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c0ff37bd590ca25063e35af745c343cb7a0271906fb7b37e4813e8f79f00268d" dependencies = [ - "bitflags", + "bitflags 2.8.0", "libc", "redox_syscall", ] @@ -2275,6 +2299,12 @@ version = "0.4.15" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d26c52dbd32dccf2d10cac7725f8eae5296885fb5703b261f7d0a0739ec807ab" +[[package]] +name = "linux-raw-sys" +version = "0.9.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fe7db12097d22ec582439daf8618b8fdd1a7bef6270e9af3b1ebcd30893cf413" + [[package]] name = "litemap" version = "0.7.4" @@ -2308,6 +2338,32 @@ dependencies = [ "scopeguard", ] +[[package]] +name = "lofty" +version = "0.22.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "781de624f162b1a8cbfbd577103ee9b8e5f62854b053ff48f4e31e68a0a7df6f" +dependencies = [ + "byteorder", + "data-encoding", + "flate2", + "lofty_attr", + "log", + "ogg_pager", + "paste", +] + +[[package]] +name = "lofty_attr" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ed9983e64b2358522f745c1251924e3ab7252d55637e80f6a0a3de642d6a9efc" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.98", +] + [[package]] name = "log" version = "0.4.25" @@ -2343,6 +2399,12 @@ dependencies = [ "digest", ] +[[package]] +name = "md5" +version = "0.7.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "490cc448043f947bae3cbee9c203358d62dbee0db12107a74be5c30ccfd09771" + [[package]] name = "memchr" version = "2.7.4" @@ -2543,6 +2605,15 @@ dependencies = [ "memchr", ] +[[package]] +name = "ogg_pager" +version = "0.7.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e034c10fb5c1c012c1b327b85df89fb0ef98ae66ec28af30f0d1eed804a40c19" +dependencies = [ + "byteorder", +] + [[package]] name = "once_cell" version = "1.20.3" @@ -2851,7 +2922,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "796d06eae7e6e74ed28ea54a8fccc584ebac84e6cf0e1e9ba41ffc807b169a01" dependencies = [ "ahash 0.8.11", - "bitflags", + "bitflags 2.8.0", "bytemuck", "chrono", "chrono-tz", @@ -2898,7 +2969,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c8e639991a8ad4fb12880ab44bcc3cf44a5703df003142334d9caf86d77d77e7" dependencies = [ "ahash 0.8.11", - "bitflags", + "bitflags 2.8.0", "hashbrown 0.15.2", "num-traits", "once_cell", @@ -2959,7 +3030,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a0a731a672dfc8ac38c1f73c9a4b2ae38d2fc8ac363bfb64c5f3a3e072ffc5ad" dependencies = [ "ahash 0.8.11", - "bitflags", + "bitflags 2.8.0", "chrono", "memchr", "once_cell", @@ -3097,7 +3168,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4f03533a93aa66127fcb909a87153a3c7cfee6f0ae59f497e73d7736208da54c" dependencies = [ "ahash 0.8.11", - "bitflags", + "bitflags 2.8.0", "bytemuck", "bytes", "chrono", @@ -3128,7 +3199,7 @@ version = "0.46.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6bf47f7409f8e75328d7d034be390842924eb276716d0458607be0bddb8cc839" dependencies = [ - "bitflags", + "bitflags 2.8.0", "bytemuck", "polars-arrow", "polars-compute", @@ -3428,7 +3499,7 @@ version = "11.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "529468c1335c1c03919960dfefdb1b3648858c20d7ec2d0663e728e4a717efbc" dependencies = [ - "bitflags", + "bitflags 2.8.0", ] [[package]] @@ -3494,7 +3565,7 @@ version = "0.5.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "03a862b389f93e68874fbf580b9de08dd02facb9a788ebadaf4a3fd33cf58834" dependencies = [ - "bitflags", + "bitflags 2.8.0", ] [[package]] @@ -3693,10 +3764,23 @@ version = "0.38.44" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fdb5bc1ae2baa591800df16c9ca78619bf65c0488b41b96ccec5d11220d8c154" dependencies = [ - "bitflags", + "bitflags 2.8.0", + "errno", + "libc", + "linux-raw-sys 0.4.15", + "windows-sys 0.59.0", +] + +[[package]] +name = "rustix" +version = "1.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e56a18552996ac8d29ecc3b190b4fdbb2d91ca4ec396de7bbffaf43f3d637e96" +dependencies = [ + "bitflags 2.8.0", "errno", "libc", - "linux-raw-sys", + "linux-raw-sys 0.9.3", "windows-sys 0.59.0", ] @@ -3795,7 +3879,7 @@ version = "2.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "897b2245f0b511c87893af39b033e5ca9cce68824c4d7e7630b5a1d339658d02" dependencies = [ - "bitflags", + "bitflags 2.8.0", "core-foundation", "core-foundation-sys", "libc", @@ -3910,6 +3994,19 @@ dependencies = [ "digest", ] +[[package]] +name = "sha256" +version = "1.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f880fc8562bdeb709793f00eb42a2ad0e672c4f883bbe59122b926eca935c8f6" +dependencies = [ + "async-trait", + "bytes", + "hex", + "sha2", + "tokio", +] + [[package]] name = "shlex" version = "1.3.0" @@ -4155,7 +4252,7 @@ checksum = "4560278f0e00ce64938540546f59f590d60beee33fffbd3b9cd47851e5fff233" dependencies = [ "atoi", "base64", - "bitflags", + "bitflags 2.8.0", "byteorder", "bytes", "chrono", @@ -4198,7 +4295,7 @@ checksum = "c5b98a57f363ed6764d5b3a12bfedf62f07aa16e1856a7ddc2a0bb190a959613" dependencies = [ "atoi", "base64", - "bitflags", + "bitflags 2.8.0", "byteorder", "chrono", "crc", @@ -4356,6 +4453,201 @@ version = "2.6.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292" +[[package]] +name = "symphonia" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "815c942ae7ee74737bb00f965fa5b5a2ac2ce7b6c01c0cc169bbeaf7abd5f5a9" +dependencies = [ + "lazy_static", + "symphonia-bundle-flac", + "symphonia-bundle-mp3", + "symphonia-codec-aac", + "symphonia-codec-adpcm", + "symphonia-codec-alac", + "symphonia-codec-pcm", + "symphonia-codec-vorbis", + "symphonia-core", + "symphonia-format-caf", + "symphonia-format-isomp4", + "symphonia-format-mkv", + "symphonia-format-ogg", + "symphonia-format-riff", + "symphonia-metadata", +] + +[[package]] +name = "symphonia-bundle-flac" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "72e34f34298a7308d4397a6c7fbf5b84c5d491231ce3dd379707ba673ab3bd97" +dependencies = [ + "log", + "symphonia-core", + "symphonia-metadata", + "symphonia-utils-xiph", +] + +[[package]] +name = "symphonia-bundle-mp3" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c01c2aae70f0f1fb096b6f0ff112a930b1fb3626178fba3ae68b09dce71706d4" +dependencies = [ + "lazy_static", + "log", + "symphonia-core", + "symphonia-metadata", +] + +[[package]] +name = "symphonia-codec-aac" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cdbf25b545ad0d3ee3e891ea643ad115aff4ca92f6aec472086b957a58522f70" +dependencies = [ + "lazy_static", + "log", + "symphonia-core", +] + +[[package]] +name = "symphonia-codec-adpcm" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c94e1feac3327cd616e973d5be69ad36b3945f16b06f19c6773fc3ac0b426a0f" +dependencies = [ + "log", + "symphonia-core", +] + +[[package]] +name = "symphonia-codec-alac" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2d8a6666649a08412906476a8b0efd9b9733e241180189e9f92b09c08d0e38f3" +dependencies = [ + "log", + "symphonia-core", +] + +[[package]] +name = "symphonia-codec-pcm" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f395a67057c2ebc5e84d7bb1be71cce1a7ba99f64e0f0f0e303a03f79116f89b" +dependencies = [ + "log", + "symphonia-core", +] + +[[package]] +name = "symphonia-codec-vorbis" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a98765fb46a0a6732b007f7e2870c2129b6f78d87db7987e6533c8f164a9f30" +dependencies = [ + "log", + "symphonia-core", + "symphonia-utils-xiph", +] + +[[package]] +name = "symphonia-core" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "798306779e3dc7d5231bd5691f5a813496dc79d3f56bf82e25789f2094e022c3" +dependencies = [ + "arrayvec", + "bitflags 1.3.2", + "bytemuck", + "lazy_static", + "log", +] + +[[package]] +name = "symphonia-format-caf" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e43c99c696a388295a29fe71b133079f5d8b18041cf734c5459c35ad9097af50" +dependencies = [ + "log", + "symphonia-core", + "symphonia-metadata", +] + +[[package]] +name = "symphonia-format-isomp4" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "abfdf178d697e50ce1e5d9b982ba1b94c47218e03ec35022d9f0e071a16dc844" +dependencies = [ + "encoding_rs", + "log", + "symphonia-core", + "symphonia-metadata", + "symphonia-utils-xiph", +] + +[[package]] +name = "symphonia-format-mkv" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1bb43471a100f7882dc9937395bd5ebee8329298e766250b15b3875652fe3d6f" +dependencies = [ + "lazy_static", + "log", + "symphonia-core", + "symphonia-metadata", + "symphonia-utils-xiph", +] + +[[package]] +name = "symphonia-format-ogg" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ada3505789516bcf00fc1157c67729eded428b455c27ca370e41f4d785bfa931" +dependencies = [ + "log", + "symphonia-core", + "symphonia-metadata", + "symphonia-utils-xiph", +] + +[[package]] +name = "symphonia-format-riff" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "05f7be232f962f937f4b7115cbe62c330929345434c834359425e043bfd15f50" +dependencies = [ + "extended", + "log", + "symphonia-core", + "symphonia-metadata", +] + +[[package]] +name = "symphonia-metadata" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bc622b9841a10089c5b18e99eb904f4341615d5aa55bbf4eedde1be721a4023c" +dependencies = [ + "encoding_rs", + "lazy_static", + "log", + "symphonia-core", +] + +[[package]] +name = "symphonia-utils-xiph" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "484472580fa49991afda5f6550ece662237b00c6f562c7d9638d1b086ed010fe" +dependencies = [ + "symphonia-core", + "symphonia-metadata", +] + [[package]] name = "syn" version = "1.0.109" @@ -4430,15 +4722,14 @@ dependencies = [ [[package]] name = "tempfile" -version = "3.17.1" +version = "3.19.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "22e5a0acb1f3f55f65cc4a866c361b2fb2a0ff6366785ae6fbb5f85df07ba230" +checksum = "7437ac7763b9b123ccf33c338a5cc1bac6f69b45a136c19bdd8a65e3916435bf" dependencies = [ - "cfg-if", "fastrand", "getrandom 0.3.1", "once_cell", - "rustix", + "rustix 1.0.3", "windows-sys 0.59.0", ] @@ -5271,7 +5562,7 @@ version = "0.33.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3268f3d866458b787f390cf61f4bbb563b922d091359f9608842999eaee3943c" dependencies = [ - "bitflags", + "bitflags 2.8.0", ] [[package]] @@ -5302,8 +5593,8 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e105d177a3871454f754b33bb0ee637ecaaac997446375fd3e5d43a2ed00c909" dependencies = [ "libc", - "linux-raw-sys", - "rustix", + "linux-raw-sys 0.4.15", + "rustix 0.38.44", ] [[package]] diff --git a/crates/dropbox/Cargo.toml b/crates/dropbox/Cargo.toml index 98f8f3b3..f21b5038 100644 --- a/crates/dropbox/Cargo.toml +++ b/crates/dropbox/Cargo.toml @@ -14,8 +14,11 @@ async-nats = "0.39.0" chrono = { version = "0.4.39", features = ["serde"] } ctr = "0.9.2" dotenv = "0.15.0" +futures = "0.3.31" hex = "0.4.3" jsonwebtoken = "9.3.1" +lofty = "0.22.2" +md5 = "0.7.0" owo-colors = "4.1.0" redis = "0.29.0" reqwest = { version = "0.12.12", features = [ @@ -26,6 +29,7 @@ reqwest = { version = "0.12.12", features = [ ], default-features = false } serde = { version = "1.0.217", features = ["derive"] } serde_json = "1.0.139" +sha256 = "1.6.0" sqlx = { version = "0.8.3", features = [ "runtime-tokio", "tls-rustls", @@ -34,5 +38,7 @@ sqlx = { version = "0.8.3", features = [ "derive", "macros", ] } +symphonia = { version = "0.5.4", features = ["all"] } +tempfile = "3.19.1" tokio = { version = "1.43.0", features = ["full"] } tokio-stream = { version = "0.1.17", features = ["full"] } diff --git a/crates/dropbox/src/consts.rs b/crates/dropbox/src/consts.rs new file mode 100644 index 00000000..e2feb2ea --- /dev/null +++ b/crates/dropbox/src/consts.rs @@ -0,0 +1,4 @@ +pub const AUDIO_EXTENSIONS: [&str; 18] = [ + "mp3", "ogg", "flac", "m4a", "aac", "mp4", "alac", "wav", "wv", "mpc", "aiff", "aif", "ac3", + "opus", "spx", "sid", "ape", "wma", +]; diff --git a/crates/dropbox/src/main.rs b/crates/dropbox/src/main.rs index 7d5783af..3ba03b65 100644 --- a/crates/dropbox/src/main.rs +++ b/crates/dropbox/src/main.rs @@ -1,9 +1,10 @@ -use std::{env, sync::Arc}; +use std::{env, sync::Arc, thread}; use actix_web::{get, post, web::{self, Data}, App, HttpRequest, HttpResponse, HttpServer, Responder}; use anyhow::Error; use dotenv::dotenv; use handlers::handle; use owo_colors::OwoColorize; +use scan::scan_dropbox; use serde_json::json; use sqlx::{postgres::PgPoolOptions, Pool, Postgres}; @@ -13,6 +14,9 @@ pub mod handlers; pub mod repo; pub mod types; pub mod client; +pub mod consts; +pub mod scan; +pub mod token; #[get("/")] async fn index(_req: HttpRequest) -> HttpResponse { @@ -49,6 +53,16 @@ async fn main() -> Result<(), Box> { let pool = PgPoolOptions::new().max_connections(5).connect(&env::var("XATA_POSTGRES_URL")?).await?; let conn = Arc::new(pool); + let cloned_conn = conn.clone(); + + thread::spawn(move || { + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on( + scan_dropbox(cloned_conn) + )?; + Ok::<(), Error>(()) + }); + let conn = conn.clone(); HttpServer::new(move || { App::new() diff --git a/crates/dropbox/src/repo/dropbox_path.rs b/crates/dropbox/src/repo/dropbox_path.rs new file mode 100644 index 00000000..596fe0cf --- /dev/null +++ b/crates/dropbox/src/repo/dropbox_path.rs @@ -0,0 +1,25 @@ +use sqlx::{Pool, Postgres}; + +use crate::xata::track::Track; +use crate::types::file::Entry; + +pub async fn create_dropbox_path( + pool: &Pool, + file: &Entry, + track: &Track, + dropbox_id: &str +) -> Result<(), sqlx::Error> { + sqlx::query(r#" + INSERT INTO dropbox_paths (dropbox_id, path, track_id, name) + VALUES ($1, $2, $3, $4) + ON CONFLICT (dropbox_id, track_id) DO NOTHING + "#) + .bind(dropbox_id) + .bind(&file.path_display) + .bind(&track.xata_id) + .bind(&file.name) + .execute(pool) + .await?; + + Ok(()) +} \ No newline at end of file diff --git a/crates/dropbox/src/repo/dropbox_token.rs b/crates/dropbox/src/repo/dropbox_token.rs index 8096fc05..7f387ecf 100644 --- a/crates/dropbox/src/repo/dropbox_token.rs +++ b/crates/dropbox/src/repo/dropbox_token.rs @@ -5,7 +5,14 @@ use crate::xata::dropbox_token::DropboxTokenWithDid; pub async fn find_dropbox_refresh_token(pool: &Pool, did: &str) -> Result, Error> { let results: Vec = sqlx::query_as(r#" - SELECT * FROM dropbox d + SELECT + d.xata_id, + d.xata_version, + d.xata_createdat, + d.xata_updatedat, + u.did as did, + dt.refresh_token as refresh_token + FROM dropbox d LEFT JOIN users u ON d.user_id = u.xata_id LEFT JOIN dropbox_tokens dt ON d.dropbox_token_id = dt.xata_id WHERE u.did = $1 @@ -20,3 +27,22 @@ pub async fn find_dropbox_refresh_token(pool: &Pool, did: &str) -> Res Ok(Some(results[0].refresh_token.clone())) } + +pub async fn find_dropbox_refresh_tokens(pool: &Pool) -> Result, Error> { + let results: Vec = sqlx::query_as(r#" + SELECT + d.xata_id, + d.xata_version, + d.xata_createdat, + d.xata_updatedat, + u.did, + dt.refresh_token as refresh_token + FROM dropbox d + LEFT JOIN users u ON d.user_id = u.xata_id + LEFT JOIN dropbox_tokens dt ON d.dropbox_token_id = dt.xata_id + "#) + .fetch_all(pool) + .await?; + + Ok(results) +} diff --git a/crates/dropbox/src/repo/mod.rs b/crates/dropbox/src/repo/mod.rs index ec7354ac..caa4331a 100644 --- a/crates/dropbox/src/repo/mod.rs +++ b/crates/dropbox/src/repo/mod.rs @@ -1 +1,3 @@ +pub mod dropbox_path; pub mod dropbox_token; +pub mod track; diff --git a/crates/dropbox/src/repo/track.rs b/crates/dropbox/src/repo/track.rs new file mode 100644 index 00000000..5380b5a5 --- /dev/null +++ b/crates/dropbox/src/repo/track.rs @@ -0,0 +1,19 @@ +use anyhow::Error; +use sqlx::{Pool, Postgres}; + +use crate::xata::track::Track; + +pub async fn get_track_by_hash(pool: &Pool, sha256: &str) -> Result, Error> { + let results: Vec = sqlx::query_as(r#" + SELECT * FROM tracks WHERE sha256 = $1 + "#) + .bind(sha256) + .fetch_all(pool) + .await?; + + if results.len() == 0 { + return Ok(None); + } + + Ok(Some(results[0].clone())) +} \ No newline at end of file diff --git a/crates/dropbox/src/scan.rs b/crates/dropbox/src/scan.rs new file mode 100644 index 00000000..320b39ea --- /dev/null +++ b/crates/dropbox/src/scan.rs @@ -0,0 +1,316 @@ +use std::{env, fs::File, io::Write, path::Path, sync::Arc}; + +use anyhow::Error; +use futures::future::BoxFuture; +use lofty::{file::TaggedFileExt, picture::{MimeType, Picture}, probe::Probe, tag::Accessor}; +use owo_colors::OwoColorize; +use reqwest::{multipart, Client}; +use serde_json::json; +use sqlx::{Pool, Postgres}; +use symphonia::core::{formats::FormatOptions, io::MediaSourceStream, meta::MetadataOptions, probe::Hint}; +use tempfile::TempDir; + +use crate::{ + client::{get_access_token, BASE_URL, CONTENT_URL}, + consts::AUDIO_EXTENSIONS, crypto::decrypt_aes_256_ctr, + repo::{dropbox_path::create_dropbox_path, dropbox_token::find_dropbox_refresh_tokens, track::get_track_by_hash}, + token::generate_token, + types::file::{Entry, EntryList} +}; + +pub async fn scan_dropbox(pool: Arc>) -> Result<(), Error>{ + let refresh_tokens = find_dropbox_refresh_tokens(&pool).await?; + for token in refresh_tokens { + let refresh_token = decrypt_aes_256_ctr( + &token.refresh_token, + &hex::decode(env::var("SPOTIFY_ENCRYPTION_KEY")?)? + )?; + + let res = get_access_token(&refresh_token).await?; + scan_audio_files( + pool.clone(), + "/Music".to_string(), + res.access_token, + token.did, + token.xata_id + ).await?; + } + Ok(()) +} + +pub fn scan_audio_files( + pool: Arc>, + path: String, + access_token: String, + did: String, + dropbox_id: String, +) -> BoxFuture<'static, Result<(), Error>> { + Box::pin(async move { + let client = Client::new(); + + let res = client.post(&format!("{}/files/get_metadata", BASE_URL)) + .bearer_auth(&access_token) + .json(&json!({ "path": path })) + .send() + .await?; + + if res.status().as_u16() == 400 || res.status().as_u16() == 409 { + println!("Path not found: {}", path.bright_red()); + return Ok(()); + } + + let entry = res.json::().await?; + + if entry.tag.clone().unwrap().as_str() == "folder" { + println!("Scanning folder: {}", path.bright_green()); + let res = client.post(&format!("{}/files/list_folder", BASE_URL)) + .bearer_auth(&access_token) + .json(&json!({ "path": path })) + .send() + .await?; + + let entries = res.json::().await?; + + for entry in entries.entries { + scan_audio_files(pool.clone(), entry.path_display, access_token.clone(), did.clone(), dropbox_id.clone()).await?; + tokio::time::sleep(std::time::Duration::from_secs(3)).await; + } + + return Ok(()); + } + + if !AUDIO_EXTENSIONS + .into_iter() + .any(|ext| path.ends_with(&format!(".{}", ext))) + { + return Ok(()); + } + + let client = Client::new(); + + println!("Downloading file: {}", path.bright_green()); + + let res = client.post(&format!("{}/files/download", CONTENT_URL)) + .bearer_auth(&access_token) + .header("Dropbox-API-Arg", &json!({ "path": path }).to_string()) + .send() + .await?; + + let bytes = res.bytes().await?; + + let temp_dir = TempDir::new()?; + let tmppath = temp_dir.path().join(&format!("{}", entry.name)); + let mut tmpfile = File::create(&tmppath)?; + tmpfile.write_all(&bytes)?; + + println!("Reading file: {}", &tmppath.clone().display().to_string().bright_green()); + + let tagged_file = match Probe::open(&tmppath)?.read() + { + Ok(tagged_file) => tagged_file, + Err(e) => { + println!("Error opening file: {}", e); + return Ok(()); + } + }; + + let primary_tag = tagged_file.primary_tag(); + let tag = match primary_tag { + Some(tag) => tag, + None => { + println!("No tag found in file"); + return Ok(()); + } + }; + + let pictures = tag.pictures(); + + println!("Title: {}", tag.get_string(&lofty::tag::ItemKey::TrackTitle).unwrap_or_default().bright_green()); + println!("Artist: {}", tag.get_string(&lofty::tag::ItemKey::TrackArtist).unwrap_or_default().bright_green()); + println!("Album Artist: {}", tag.get_string(&lofty::tag::ItemKey::AlbumArtist).unwrap_or_default().bright_green()); + println!("Album: {}", tag.get_string(&lofty::tag::ItemKey::AlbumTitle).unwrap_or_default().bright_green()); + println!("Lyrics: {}", tag.get_string(&lofty::tag::ItemKey::Lyrics).unwrap_or_default().bright_green()); + println!("Year: {}", tag.year().unwrap_or_default().bright_green()); + println!("Track Number: {}", tag.track().unwrap_or_default().bright_green()); + println!("Track Total: {}", tag.track_total().unwrap_or_default().bright_green()); + println!("Release Date: {:?}", tag.get_string(&lofty::tag::ItemKey::OriginalReleaseDate).unwrap_or_default().bright_green()); + println!("Recording Date: {:?}", tag.get_string(&lofty::tag::ItemKey::RecordingDate).unwrap_or_default().bright_green()); + println!("Copyright Message: {}", tag.get_string(&lofty::tag::ItemKey::CopyrightMessage).unwrap_or_default().bright_green()); + println!("Pictures: {:?}", pictures); + + let title = tag.get_string(&lofty::tag::ItemKey::TrackTitle).unwrap_or_default(); + let artist = tag.get_string(&lofty::tag::ItemKey::TrackArtist).unwrap_or_default(); + let album = tag.get_string(&lofty::tag::ItemKey::AlbumTitle).unwrap_or_default(); + let album_artist = tag.get_string(&lofty::tag::ItemKey::AlbumArtist).unwrap_or_default(); + + let access_token = generate_token(&did)?; + + // check if track exists + // + // if not, create track + // upload album art + // + // link path to track + + let hash = sha256::digest( + format!("{} - {} - {}", title, artist, album).to_lowercase(), + ); + + let track = get_track_by_hash(&pool, &hash).await?; + let duration = get_track_duration(&tmppath).await?; + let albumart_id = md5::compute(&format!("{} - {}", album_artist, album).to_lowercase()); + let albumart_id = format!("{:x}", albumart_id); + + match track { + Some(track) => { + println!("Track exists: {}", title.bright_green()); + let status = create_dropbox_path( + &pool, + &entry, + &track, + &dropbox_id, + ) + .await; + println!("status {:?}", status); + }, + None => { + println!("Creating track: {}", title.bright_green()); + let album_art = upload_album_cover(albumart_id.into(), pictures, &access_token).await?; + let client = Client::new(); + const URL: &str = "https://api.rocksky.app/tracks"; + let response = client + .post(URL) + .header("Authorization", format!("Bearer {}", access_token)) + .json(&serde_json::json!({ + "title": tag.get_string(&lofty::tag::ItemKey::TrackTitle), + "album": tag.get_string(&lofty::tag::ItemKey::AlbumTitle), + "artist": tag.get_string(&lofty::tag::ItemKey::TrackArtist), + "albumArtist": tag.get_string(&lofty::tag::ItemKey::AlbumArtist), + "duration": duration, + "trackNumber": tag.track(), + "releaseDate": tag.get_string(&lofty::tag::ItemKey::OriginalReleaseDate).map(|date| match date.contains("-") { + true => Some(date), + false => None, + }), + "year": tag.year(), + "discNumber": tag.disk(), + "composer": tag.get_string(&lofty::tag::ItemKey::Composer), + "albumArt": match album_art{ + Some(album_art) => Some(format!("https://cdn.rocksky.app/covers/{}", album_art)), + None => None + }, + "lyrics": tag.get_string(&lofty::tag::ItemKey::Lyrics), + "copyrightMessage": tag.get_string(&lofty::tag::ItemKey::CopyrightMessage), + })) + .send() + .await?; + println!("Track Saved: {} {}", title, response.status()); + + + let track = get_track_by_hash(&pool, &hash).await?; + if let Some(track) = track { + create_dropbox_path( + &pool, + &entry, + &track, + &dropbox_id, + ) + .await?; + return Ok(()); + } + + println!("Failed to create track: {}", title.bright_green()); + } + } + + Ok(()) + }) +} + +pub async fn upload_album_cover(name: String, pictures: &[Picture], token: &str) -> Result, Error> { + if pictures.is_empty() { + return Ok(None); + } + + let picture = &pictures[0]; + + let buffer = match picture.mime_type() { + Some(MimeType::Jpeg) => Some(picture.data().to_vec()), + Some(MimeType::Png) => Some(picture.data().to_vec()), + Some(MimeType::Gif) => Some(picture.data().to_vec()), + Some(MimeType::Bmp) => Some(picture.data().to_vec()), + Some(MimeType::Tiff) => Some(picture.data().to_vec()), + _ => None + }; + + if buffer.is_none() { + return Ok(None); + } + + let buffer = buffer.unwrap(); + + let ext = match picture.mime_type() { + Some(MimeType::Jpeg) => "jpg", + Some(MimeType::Png) => "png", + Some(MimeType::Gif) => "gif", + Some(MimeType::Bmp) => "bmp", + Some(MimeType::Tiff) => "tiff", + _ => { + return Ok(None); + } + }; + + let name = format!("{}.{}", name, ext); + + let part = multipart::Part::bytes(buffer).file_name(name.clone()); + let form = multipart::Form::new().part("file", part); + let client = Client::new(); + + const URL: &str = "https://uploads.rocksky.app"; + + let response = client + .post(URL) + .header("Authorization", format!("Bearer {}", token)) + .multipart(form) + .send() + .await?; + + println!("Cover uploaded: {}", response.status()); + + Ok(Some(name)) +} + + + +pub async fn get_track_duration(path: &Path) -> Result { + let duration = 0; + let media_source = MediaSourceStream::new(Box::new(std::fs::File::open(path)?), Default::default()); + let mut hint = Hint::new(); + + if let Some(extension) = path.extension() { + if let Some(extension) = extension.to_str() { + hint.with_extension(extension); + } + } + + + let meta_opts = MetadataOptions::default(); + let format_opts = FormatOptions::default(); + + let probed = match symphonia::default::get_probe().format(&hint, media_source, &format_opts, &meta_opts) { + Ok(probed) => probed, + Err(_) => { + println!("Error probing file"); + return Ok(duration); + }, + }; + + if let Some(track) = probed.format.tracks().first() { + if let Some(duration) = track.codec_params.n_frames { + if let Some(sample_rate) = track.codec_params.sample_rate { + return Ok((duration as f64 / sample_rate as f64) as u64 * 1000); + } + } +} + Ok(duration) +} \ No newline at end of file diff --git a/crates/dropbox/src/token.rs b/crates/dropbox/src/token.rs new file mode 100644 index 00000000..40c8566c --- /dev/null +++ b/crates/dropbox/src/token.rs @@ -0,0 +1,64 @@ +use std::env; + +use anyhow::Error; +use jsonwebtoken::DecodingKey; +use jsonwebtoken::EncodingKey; +use jsonwebtoken::Header; +use jsonwebtoken::Validation; +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Serialize, Deserialize)] +pub struct Claims { + exp: usize, + iat: usize, + did: String, +} + +pub fn generate_token(did: &str) -> Result { + if env::var("JWT_SECRET").is_err() { + return Err(Error::msg("JWT_SECRET is not set")); + } + + let claims = Claims { + exp: chrono::Utc::now().timestamp() as usize + 3600, + iat: chrono::Utc::now().timestamp() as usize, + did: did.to_string(), + }; + + jsonwebtoken::encode( + &Header::default(), + &claims, + &EncodingKey::from_secret(env::var("JWT_SECRET")?.as_ref()), + ) + .map_err(Into::into) +} + +pub fn decode_token(token: &str) -> Result { + if env::var("JWT_SECRET").is_err() { + return Err(Error::msg("JWT_SECRET is not set")); + } + + jsonwebtoken::decode::( + token, + &DecodingKey::from_secret(env::var("JWT_SECRET")?.as_ref()), + &Validation::default(), + ) + .map(|data| data.claims) + .map_err(Into::into) +} + +#[cfg(test)] +mod tests { + use dotenv::dotenv; + + use super::*; + + #[test] + fn test_generate_token() { + dotenv().ok(); + let token = generate_token("did:plc:7vdlgi2bflelz7mmuxoqjfcr").unwrap(); + let claims = decode_token(&token).unwrap(); + + assert_eq!(claims.did, "did:plc:7vdlgi2bflelz7mmuxoqjfcr"); + } +} diff --git a/crates/googledrive/Cargo.toml b/crates/googledrive/Cargo.toml index 9f2e1a8f..3675526c 100644 --- a/crates/googledrive/Cargo.toml +++ b/crates/googledrive/Cargo.toml @@ -14,8 +14,11 @@ async-nats = "0.39.0" chrono = { version = "0.4.39", features = ["serde"] } ctr = "0.9.2" dotenv = "0.15.0" +futures = "0.3.31" hex = "0.4.3" jsonwebtoken = "9.3.1" +lofty = "0.22.2" +md5 = "0.7.0" owo-colors = "4.1.0" redis = "0.29.0" reqwest = { version = "0.12.12", features = [ @@ -27,6 +30,7 @@ reqwest = { version = "0.12.12", features = [ serde = { version = "1.0.217", features = ["derive"] } serde_json = "1.0.139" serde_urlencoded = "0.7.1" +sha256 = "1.6.0" sqlx = { version = "0.8.3", features = [ "runtime-tokio", "tls-rustls", @@ -35,5 +39,7 @@ sqlx = { version = "0.8.3", features = [ "derive", "macros", ] } +symphonia = { version = "0.5.4", features = ["all"] } +tempfile = "3.19.1" tokio = { version = "1.43.0", features = ["full"] } tokio-stream = { version = "0.1.17", features = ["full"] } diff --git a/crates/googledrive/src/consts.rs b/crates/googledrive/src/consts.rs new file mode 100644 index 00000000..e2feb2ea --- /dev/null +++ b/crates/googledrive/src/consts.rs @@ -0,0 +1,4 @@ +pub const AUDIO_EXTENSIONS: [&str; 18] = [ + "mp3", "ogg", "flac", "m4a", "aac", "mp4", "alac", "wav", "wv", "mpc", "aiff", "aif", "ac3", + "opus", "spx", "sid", "ape", "wma", +]; diff --git a/crates/googledrive/src/main.rs b/crates/googledrive/src/main.rs index fcf68f49..f4d9edef 100644 --- a/crates/googledrive/src/main.rs +++ b/crates/googledrive/src/main.rs @@ -1,18 +1,22 @@ -use std::{env, sync::Arc}; +use std::{env, sync::Arc, thread}; use actix_web::{get, post, web::{self, Data}, App, HttpRequest, HttpResponse, HttpServer, Responder}; use anyhow::Error; use dotenv::dotenv; use handlers::handle; use owo_colors::OwoColorize; +use scan::scan_googledrive; use serde_json::json; use sqlx::{postgres::PgPoolOptions, Pool, Postgres}; +pub mod token; pub mod xata; pub mod crypto; pub mod handlers; pub mod repo; pub mod types; pub mod client; +pub mod consts; +pub mod scan; #[get("/")] async fn index(_req: HttpRequest) -> HttpResponse { @@ -50,6 +54,16 @@ async fn main() -> Result<(), Box> { let pool = PgPoolOptions::new().max_connections(5).connect(&env::var("XATA_POSTGRES_URL")?).await?; let conn = Arc::new(pool); + let cloned_conn = conn.clone(); + + thread::spawn(move || { + let rt = tokio::runtime::Runtime::new().unwrap(); + rt.block_on( + scan_googledrive(cloned_conn) + )?; + Ok::<(), Error>(()) + }); + let conn = conn.clone(); HttpServer::new(move || { App::new() diff --git a/crates/googledrive/src/repo/google_drive_path.rs b/crates/googledrive/src/repo/google_drive_path.rs new file mode 100644 index 00000000..518737b1 --- /dev/null +++ b/crates/googledrive/src/repo/google_drive_path.rs @@ -0,0 +1,24 @@ +use sqlx::{Pool, Postgres}; + +use crate::{types::file::File, xata::track::Track}; + +pub async fn create_google_drive_path( + pool: &Pool, + file: &File, + track: &Track, + google_drive_id: &str +) -> Result<(), sqlx::Error> { + sqlx::query(r#" + INSERT INTO google_drive_paths (google_drive_id, file_id, track_id, name) + VALUES ($1, $2, $3, $4) + ON CONFLICT (google_drive_id, file_id, track_id) DO NOTHING + "#) + .bind(google_drive_id) + .bind(&file.id) + .bind(&track.xata_id) + .bind(&file.name) + .execute(pool) + .await?; + + Ok(()) +} \ No newline at end of file diff --git a/crates/googledrive/src/repo/google_drive_token.rs b/crates/googledrive/src/repo/google_drive_token.rs index b7024e4c..2268ec5f 100644 --- a/crates/googledrive/src/repo/google_drive_token.rs +++ b/crates/googledrive/src/repo/google_drive_token.rs @@ -5,7 +5,14 @@ use crate::xata::google_drive_token::GoogleDriveTokenWithDid; pub async fn find_google_drive_refresh_token(pool: &Pool, did: &str) -> Result, Error> { let results: Vec = sqlx::query_as(r#" - SELECT * FROM google_drive gd + SELECT + gd.xata_id, + gd.xata_version, + gd.xata_createdat, + gd.xata_updatedat, + u.did, + gt.refresh_token + FROM google_drive gd LEFT JOIN users u ON gd.user_id = u.xata_id LEFT JOIN google_drive_tokens gt ON gd.google_drive_token_id = gt.xata_id WHERE u.did = $1 @@ -20,3 +27,22 @@ pub async fn find_google_drive_refresh_token(pool: &Pool, did: &str) - Ok(Some(results[0].refresh_token.clone())) } + +pub async fn find_google_drive_refresh_tokens(pool: &Pool) -> Result, Error> { + let results: Vec = sqlx::query_as(r#" + SELECT + gd.xata_id, + gd.xata_version, + gd.xata_createdat, + gd.xata_updatedat, + u.did, + gt.refresh_token + FROM google_drive gd + LEFT JOIN users u ON gd.user_id = u.xata_id + LEFT JOIN google_drive_tokens gt ON gd.google_drive_token_id = gt.xata_id + "#) + .fetch_all(pool) + .await?; + + Ok(results) +} diff --git a/crates/googledrive/src/repo/mod.rs b/crates/googledrive/src/repo/mod.rs index 3181c60e..9abf23a7 100644 --- a/crates/googledrive/src/repo/mod.rs +++ b/crates/googledrive/src/repo/mod.rs @@ -1 +1,3 @@ +pub mod google_drive_path; pub mod google_drive_token; +pub mod track; diff --git a/crates/googledrive/src/repo/track.rs b/crates/googledrive/src/repo/track.rs new file mode 100644 index 00000000..5380b5a5 --- /dev/null +++ b/crates/googledrive/src/repo/track.rs @@ -0,0 +1,19 @@ +use anyhow::Error; +use sqlx::{Pool, Postgres}; + +use crate::xata::track::Track; + +pub async fn get_track_by_hash(pool: &Pool, sha256: &str) -> Result, Error> { + let results: Vec = sqlx::query_as(r#" + SELECT * FROM tracks WHERE sha256 = $1 + "#) + .bind(sha256) + .fetch_all(pool) + .await?; + + if results.len() == 0 { + return Ok(None); + } + + Ok(Some(results[0].clone())) +} \ No newline at end of file diff --git a/crates/googledrive/src/scan.rs b/crates/googledrive/src/scan.rs new file mode 100644 index 00000000..2886b63e --- /dev/null +++ b/crates/googledrive/src/scan.rs @@ -0,0 +1,322 @@ +use std::{env, io::Write, path::Path, sync::Arc}; + +use anyhow::Error; +use futures::future::BoxFuture; +use lofty::{file::TaggedFileExt, picture::{MimeType, Picture}, probe::Probe, tag::Accessor}; +use owo_colors::OwoColorize; +use reqwest::{multipart, Client}; +use sqlx::{Pool, Postgres}; +use symphonia::core::{formats::FormatOptions, io::MediaSourceStream, meta::MetadataOptions, probe::Hint}; +use tempfile::TempDir; + +use crate::{client::{GoogleDriveClient, BASE_URL}, consts::AUDIO_EXTENSIONS, crypto::decrypt_aes_256_ctr, repo::{google_drive_path::create_google_drive_path, google_drive_token::find_google_drive_refresh_tokens, track::get_track_by_hash}, token::generate_token, types::file::{File, FileList}}; + +pub async fn scan_googledrive(pool: Arc>) -> Result<(), Error> { + let refresh_tokens = find_google_drive_refresh_tokens(&pool).await?; + for token in refresh_tokens { + let refresh_token = decrypt_aes_256_ctr( + &token.refresh_token, + &hex::decode(env::var("SPOTIFY_ENCRYPTION_KEY")?)? + )?; + + let client = GoogleDriveClient::new(&refresh_token).await?; + let filelist = client.get_music_directory().await?; + let music_dir = filelist.files.first().unwrap(); + let access_token = client.access_token.clone(); + scan_audio_files( + pool.clone(), + music_dir.id.clone(), + access_token, + token.did.clone(), + token.xata_id.clone() + ).await?; + } + Ok(()) +} + +pub fn scan_audio_files( + pool: Arc>, + file_id: String, + access_token: String, + did: String, + google_drive_id: String, +) -> BoxFuture<'static, Result<(), Error>> { + Box::pin(async move { + let client = Client::new(); + let url = format!("{}/files/{}", BASE_URL, file_id); + let res = client.get(&url) + .bearer_auth(&access_token) + .query(&[ + ("fields", "id, name, mimeType, parents"), + ]) + .send() + .await?; + + let file = res.json::().await?; + + if file.mime_type == "application/vnd.google-apps.folder" { + println!("Scanning folder: {}", file.name.bright_green()); + + let url = format!("{}/files", BASE_URL); + let res = client.get(&url) + .bearer_auth(&access_token) + .query(&[ + ("q", format!("'{}' in parents", file.id).as_str()), + ("fields", "files(id, name, mimeType, parents)"), + ("orderBy", "name"), + ]) + .send() + .await?; + let filelist = res.json::().await?; + + for file in filelist.files { + scan_audio_files( + pool.clone(), + file.id, + access_token.clone(), + did.clone(), + google_drive_id.clone() + ).await?; + tokio::time::sleep(std::time::Duration::from_secs(3)).await; + } + + return Ok(()); + } + + if !AUDIO_EXTENSIONS + .into_iter() + .any(|ext| file.name.ends_with(&format!(".{}", ext))) + { + return Ok(()); + } + + println!("Downloading file: {}", file.name.bright_green()); + + let client = Client::new(); + + let url = format!("{}/files/{}", BASE_URL, file_id); + let res = client.get(&url) + .bearer_auth(&access_token) + .query(&[ + ("alt", "media"), + ]) + .send() + .await?; + + let bytes = res.bytes().await?; + + let temp_dir = TempDir::new()?; + let tmppath = temp_dir.path().join(&format!("{}", file.name)); + let mut tmpfile = std::fs::File::create(&tmppath)?; + tmpfile.write_all(&bytes)?; + + println!("Reading file: {}", &tmppath.clone().display().to_string().bright_green()); + + let tagged_file = match Probe::open(&tmppath)?.read() + { + Ok(tagged_file) => tagged_file, + Err(e) => { + println!("Error opening file: {}", e); + return Ok(()); + } + }; + + let primary_tag = tagged_file.primary_tag(); + let tag = match primary_tag { + Some(tag) => tag, + None => { + println!("No tag found in file"); + return Ok(()); + } + }; + + let pictures = tag.pictures(); + + println!("Title: {}", tag.get_string(&lofty::tag::ItemKey::TrackTitle).unwrap_or_default().bright_green()); + println!("Artist: {}", tag.get_string(&lofty::tag::ItemKey::TrackArtist).unwrap_or_default().bright_green()); + println!("Album Artist: {}", tag.get_string(&lofty::tag::ItemKey::AlbumArtist).unwrap_or_default().bright_green()); + println!("Album: {}", tag.get_string(&lofty::tag::ItemKey::AlbumTitle).unwrap_or_default().bright_green()); + println!("Lyrics: {}", tag.get_string(&lofty::tag::ItemKey::Lyrics).unwrap_or_default().bright_green()); + println!("Year: {}", tag.year().unwrap_or_default().bright_green()); + println!("Track Number: {}", tag.track().unwrap_or_default().bright_green()); + println!("Track Total: {}", tag.track_total().unwrap_or_default().bright_green()); + println!("Release Date: {:?}", tag.get_string(&lofty::tag::ItemKey::OriginalReleaseDate).unwrap_or_default().bright_green()); + println!("Recording Date: {:?}", tag.get_string(&lofty::tag::ItemKey::RecordingDate).unwrap_or_default().bright_green()); + println!("Copyright Message: {}", tag.get_string(&lofty::tag::ItemKey::CopyrightMessage).unwrap_or_default().bright_green()); + println!("Pictures: {:?}", pictures); + + let title = tag.get_string(&lofty::tag::ItemKey::TrackTitle).unwrap_or_default(); + let artist = tag.get_string(&lofty::tag::ItemKey::TrackArtist).unwrap_or_default(); + let album_artist = tag.get_string(&lofty::tag::ItemKey::AlbumArtist).unwrap_or_default(); + let album = tag.get_string(&lofty::tag::ItemKey::AlbumTitle).unwrap_or_default(); + let access_token = generate_token(&did)?; + + // check if track exists + // + // if not, create track + // upload album artist + // + // link path to track + + let hash = sha256::digest( + format!("{} - {} - {}", title, artist, album).to_lowercase(), + ); + + let track = get_track_by_hash(&pool, &hash).await?; + let duration = get_track_duration(&tmppath).await?; + let albumart_id = md5::compute(&format!("{} - {}", album_artist, album).to_lowercase()); + let albumart_id = format!("{:x}", albumart_id); + + match track { + Some(track) => { + println!("Track exists: {}", title.bright_green()); + create_google_drive_path( + &pool, + &file, + &track, + &google_drive_id, + ) + .await?; + }, + None => { + println!("Creating track: {}", title.bright_green()); + + let albumart = upload_album_cover(albumart_id.into(), pictures, &access_token).await?; + + let client = Client::new(); + const URL: &str = "https://api.rocksky.app/tracks"; + let response = client + .post(URL) + .header("Authorization", format!("Bearer {}", access_token)) + .json(&serde_json::json!({ + "title": tag.get_string(&lofty::tag::ItemKey::TrackTitle), + "album": tag.get_string(&lofty::tag::ItemKey::AlbumTitle), + "artist": tag.get_string(&lofty::tag::ItemKey::TrackArtist), + "albumArtist": tag.get_string(&lofty::tag::ItemKey::AlbumArtist), + "duration": duration, + "trackNumber": tag.track(), + "releaseDate": tag.get_string(&lofty::tag::ItemKey::OriginalReleaseDate).map(|date| match date.contains("-") { + true => Some(date), + false => None, + }), + "year": tag.year(), + "discNumber": tag.disk(), + "composer": tag.get_string(&lofty::tag::ItemKey::Composer), + "albumArt": match albumart { + Some(albumart) => Some(format!("https://cdn.rocksky.app/covers/{}", albumart)), + None => None + }, + "lyrics": tag.get_string(&lofty::tag::ItemKey::Lyrics), + "copyrightMessage": tag.get_string(&lofty::tag::ItemKey::CopyrightMessage), + })) + .send() + .await?; + println!("Track Saved: {} {}", title, response.status()); + + + let track = get_track_by_hash(&pool, &hash).await?; + if let Some(track) = track { + create_google_drive_path( + &pool, + &file, + &track, + &google_drive_id, + ) + .await?; + return Ok(()); + } + + println!("Failed to create track: {}", title.bright_green()); + } + } + + Ok(()) + }) +} + +pub async fn upload_album_cover(name: String, pictures: &[Picture], token: &str) -> Result, Error> { + if pictures.is_empty() { + return Ok(None); + } + + let picture = &pictures[0]; + + let buffer = match picture.mime_type() { + Some(MimeType::Jpeg) => Some(picture.data().to_vec()), + Some(MimeType::Png) => Some(picture.data().to_vec()), + Some(MimeType::Gif) => Some(picture.data().to_vec()), + Some(MimeType::Bmp) => Some(picture.data().to_vec()), + Some(MimeType::Tiff) => Some(picture.data().to_vec()), + _ => None + }; + + if buffer.is_none() { + return Ok(None); + } + + let buffer = buffer.unwrap(); + + + let ext = match picture.mime_type() { + Some(MimeType::Jpeg) => "jpg", + Some(MimeType::Png) => "png", + Some(MimeType::Gif) => "gif", + Some(MimeType::Bmp) => "bmp", + Some(MimeType::Tiff) => "tiff", + _ => { + return Ok(None); + } + }; + + let name = format!("{}.{}", name, ext); + + let part = multipart::Part::bytes(buffer).file_name(name.clone()); + let form = multipart::Form::new().part("file", part); + let client = Client::new(); + + const URL: &str = "https://uploads.rocksky.app"; + + let response = client + .post(URL) + .header("Authorization", format!("Bearer {}", token)) + .multipart(form) + .send() + .await?; + + println!("Cover uploaded: {}", response.status()); + + Ok(Some(name)) +} + +pub async fn get_track_duration(path: &Path) -> Result { + let duration = 0; + let media_source = MediaSourceStream::new(Box::new(std::fs::File::open(path)?), Default::default()); + let mut hint = Hint::new(); + + if let Some(extension) = path.extension() { + if let Some(extension) = extension.to_str() { + hint.with_extension(extension); + } + } + + + let meta_opts = MetadataOptions::default(); + let format_opts = FormatOptions::default(); + + let probed = match symphonia::default::get_probe().format(&hint, media_source, &format_opts, &meta_opts) { + Ok(probed) => probed, + Err(_) => { + println!("Error probing file"); + return Ok(duration); + }, + }; + + if let Some(track) = probed.format.tracks().first() { + if let Some(duration) = track.codec_params.n_frames { + if let Some(sample_rate) = track.codec_params.sample_rate { + return Ok((duration as f64 / sample_rate as f64) as u64 * 1000); + } + } +} + Ok(duration) +} \ No newline at end of file diff --git a/crates/googledrive/src/token.rs b/crates/googledrive/src/token.rs new file mode 100644 index 00000000..40c8566c --- /dev/null +++ b/crates/googledrive/src/token.rs @@ -0,0 +1,64 @@ +use std::env; + +use anyhow::Error; +use jsonwebtoken::DecodingKey; +use jsonwebtoken::EncodingKey; +use jsonwebtoken::Header; +use jsonwebtoken::Validation; +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Serialize, Deserialize)] +pub struct Claims { + exp: usize, + iat: usize, + did: String, +} + +pub fn generate_token(did: &str) -> Result { + if env::var("JWT_SECRET").is_err() { + return Err(Error::msg("JWT_SECRET is not set")); + } + + let claims = Claims { + exp: chrono::Utc::now().timestamp() as usize + 3600, + iat: chrono::Utc::now().timestamp() as usize, + did: did.to_string(), + }; + + jsonwebtoken::encode( + &Header::default(), + &claims, + &EncodingKey::from_secret(env::var("JWT_SECRET")?.as_ref()), + ) + .map_err(Into::into) +} + +pub fn decode_token(token: &str) -> Result { + if env::var("JWT_SECRET").is_err() { + return Err(Error::msg("JWT_SECRET is not set")); + } + + jsonwebtoken::decode::( + token, + &DecodingKey::from_secret(env::var("JWT_SECRET")?.as_ref()), + &Validation::default(), + ) + .map(|data| data.claims) + .map_err(Into::into) +} + +#[cfg(test)] +mod tests { + use dotenv::dotenv; + + use super::*; + + #[test] + fn test_generate_token() { + dotenv().ok(); + let token = generate_token("did:plc:7vdlgi2bflelz7mmuxoqjfcr").unwrap(); + let claims = decode_token(&token).unwrap(); + + assert_eq!(claims.did, "did:plc:7vdlgi2bflelz7mmuxoqjfcr"); + } +}