diff --git a/.gitignore b/.gitignore index 2e31e038..50477f5f 100644 --- a/.gitignore +++ b/.gitignore @@ -8,4 +8,5 @@ node_modules/ .vscode/ *.DS_Store rocksky-backup.sql -*.db \ No newline at end of file +*.db +*.parquet \ No newline at end of file diff --git a/Cargo.lock b/Cargo.lock index 5967a6ef..4019cc08 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1373,6 +1373,17 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "cron" +version = "0.15.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5877d3fbf742507b66bc2a1945106bd30dd8504019d596901ddd012a4dd01740" +dependencies = [ + "chrono", + "once_cell", + "winnow 0.6.26", +] + [[package]] name = "crossbeam-channel" version = "0.5.15" @@ -4846,6 +4857,7 @@ dependencies = [ "async-nats", "chrono", "clap", + "cron", "dotenv", "duckdb", "owo-colors", @@ -6776,7 +6788,7 @@ dependencies = [ "toml_datetime 0.7.0", "toml_parser", "toml_writer", - "winnow", + "winnow 0.7.10", ] [[package]] @@ -6808,7 +6820,7 @@ dependencies = [ "serde_spanned 0.6.9", "toml_datetime 0.6.11", "toml_write", - "winnow", + "winnow 0.7.10", ] [[package]] @@ -6817,7 +6829,7 @@ version = "1.0.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "97200572db069e74c512a14117b296ba0a80a30123fbbb5aa1f4a348f639ca30" dependencies = [ - "winnow", + "winnow 0.7.10", ] [[package]] @@ -7743,6 +7755,15 @@ version = "0.53.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "271414315aff87387382ec3d271b52d7ae78726f5d44ac98b4f4030c91880486" +[[package]] +name = "winnow" +version = "0.6.26" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e90edd2ac1aa278a5c4599b1d89cf03074b610800f866d4026dc199d7929a28" +dependencies = [ + "memchr", +] + [[package]] name = "winnow" version = "0.7.10" diff --git a/crates/analytics/Cargo.toml b/crates/analytics/Cargo.toml index a81d1436..a691ee83 100644 --- a/crates/analytics/Cargo.toml +++ b/crates/analytics/Cargo.toml @@ -29,3 +29,4 @@ clap = "4.5.31" actix-web = "4.9.0" tokio-stream = { version = "0.1.17", features = ["full"] } tracing = "0.1.41" +cron = "0.15.0" diff --git a/crates/analytics/src/lib.rs b/crates/analytics/src/lib.rs index 554f9d72..479346e9 100644 --- a/crates/analytics/src/lib.rs +++ b/crates/analytics/src/lib.rs @@ -1,6 +1,8 @@ use std::{ env, + str::FromStr, sync::{Arc, Mutex}, + thread, }; use anyhow::Error; @@ -22,6 +24,7 @@ pub async fn serve() -> Result<(), Error> { create_tables(&conn).await?; let conn = Arc::new(Mutex::new(conn)); + export_parquets(conn.clone()); cmd::serve::serve(conn).await?; Ok(()) @@ -42,3 +45,103 @@ pub async fn sync() -> Result<(), Error> { Ok(()) } + +fn export_parquets(conn: Arc>) { + thread::spawn(move || { + // fire every 1 minute + let cron_expr = "0 * * * * * *"; + let schedule = cron::Schedule::from_str(cron_expr); + if let Err(err) = schedule { + tracing::error!("Failed to parse cron expression: {}", cron_expr); + tracing::error!(error = %err); + return Ok(()); + } + let schedule = schedule.unwrap(); + loop { + let now = chrono::Utc::now(); + let mut upcoming = schedule.upcoming(chrono::Utc).take(1); + + if let Some(next) = upcoming.next() { + let duration = next.signed_duration_since(now).to_std().unwrap(); + thread::sleep(duration); + tracing::info!("Exporting parquets ..."); + + let conn = conn.lock().unwrap(); + conn.execute_batch( + "BEGIN; + COPY (SELECT * FROM scrobbles) TO 'scrobbles.parquet' (FORMAT PARQUET); + COPY (SELECT * FROM artists) TO 'artists.parquet' (FORMAT PARQUET); + COPY (SELECT * FROM albums) TO 'albums.parquet' (FORMAT PARQUET); + COPY (SELECT * FROM tracks) TO 'tracks.parquet' (FORMAT PARQUET); + COPY (SELECT * FROM users) TO 'users.parquet' (FORMAT PARQUET); + COPY (SELECT * FROM album_tracks) TO 'album_tracks.parquet' (FORMAT PARQUET); + COPY (SELECT * FROM artist_albums) TO 'artist_albums.parquet' (FORMAT PARQUET); + COPY (SELECT * FROM artist_tracks) TO 'artist_tracks.parquet' (FORMAT PARQUET); + COPY (SELECT * FROM loved_tracks) TO 'loved_tracks.parquet' (FORMAT PARQUET); + COPY (SELECT * FROM user_albums) TO 'user_albums.parquet' (FORMAT PARQUET); + COPY (SELECT * FROM user_artists) TO 'user_artists.parquet' (FORMAT PARQUET); + COPY (SELECT * FROM user_tracks) TO 'user_tracks.parquet' (FORMAT PARQUET); + COMMIT;", + )?; + + drop(conn); + + if env::var("CF_ACCOUNT_ID").is_err() { + tracing::warn!("CF_ACCOUNT_ID is not set, skipping upload to R2"); + continue; + } + + upload_to_r2("scrobbles.parquet"); + upload_to_r2("artists.parquet"); + upload_to_r2("albums.parquet"); + upload_to_r2("tracks.parquet"); + upload_to_r2("users.parquet"); + upload_to_r2("album_tracks.parquet"); + upload_to_r2("artist_albums.parquet"); + upload_to_r2("artist_tracks.parquet"); + upload_to_r2("loved_tracks.parquet"); + upload_to_r2("user_albums.parquet"); + upload_to_r2("user_artists.parquet"); + upload_to_r2("user_tracks.parquet"); + + tracing::info!("Exported parquets successfully."); + } + } + + #[allow(unreachable_code)] + Ok::<(), Error>(()) + }); +} + +fn upload_to_r2(file: &str) { + let status = std::process::Command::new("aws") + .arg("s3") + .arg("cp") + .arg(file) + .arg(format!( + "s3://{}", + env::var("R2_BUCKET_NAME").unwrap_or("rocksky-backup".to_string()) + )) + .arg("--endpoint-url") + .arg(&format!( + "https://{}.r2.cloudflarestorage.com", + env::var("CF_ACCOUNT_ID").unwrap() + )) + .arg("--profile") + .arg("r2") + .stdout(std::process::Stdio::inherit()) + .stderr(std::process::Stdio::inherit()) + .status(); + match status { + Ok(status) => { + if status.success() { + tracing::info!("Uploaded {} to R2 successfully.", file); + } else { + tracing::error!("Failed to upload {} to R2.", file); + } + } + Err(err) => { + tracing::error!("Failed to execute aws command: {}", err); + } + }; +}