diff --git a/crates/analytics/src/main.rs b/crates/analytics/src/main.rs deleted file mode 100644 index 61af8c51..00000000 --- a/crates/analytics/src/main.rs +++ /dev/null @@ -1,50 +0,0 @@ -use core::create_tables; -use std::{ - env, - sync::{Arc, Mutex}, -}; - -use clap::Command; -use cmd::{serve::serve, sync::sync}; -use dotenv::dotenv; -use duckdb::Connection; -use sqlx::postgres::PgPoolOptions; - -pub mod cmd; -pub mod core; -pub mod handlers; -pub mod subscriber; -pub mod types; -pub mod xata; - -fn cli() -> Command { - Command::new("analytics") - .version(env!("CARGO_PKG_VERSION")) - .about("Rocksky Analytics CLI built with Rust and DuckDB") - .subcommand(Command::new("sync").about("Sync data from Xata to DuckDB")) - .subcommand(Command::new("serve").about("Serve the Rocksky Analytics API")) -} - -#[tokio::main] -async fn main() -> Result<(), Box> { - dotenv().ok(); - - let pool = PgPoolOptions::new() - .max_connections(5) - .connect(&env::var("XATA_POSTGRES_URL")?) - .await?; - let conn = Connection::open("./rocksky-analytics.ddb")?; - - create_tables(&conn).await?; - - let args = cli().get_matches(); - let conn = Arc::new(Mutex::new(conn)); - - match args.subcommand() { - Some(("sync", _)) => sync(conn, &pool).await?, - Some(("serve", _)) => serve(conn).await?, - _ => serve(conn).await?, - } - - Ok(()) -} diff --git a/crates/dropbox/src/main.rs b/crates/dropbox/src/main.rs deleted file mode 100644 index 7f9d0fd6..00000000 --- a/crates/dropbox/src/main.rs +++ /dev/null @@ -1,37 +0,0 @@ -use clap::Command; -use cmd::{scan::scan, serve::serve}; -use dotenv::dotenv; - -pub mod client; -pub mod cmd; -pub mod consts; -pub mod crypto; -pub mod handlers; -pub mod repo; -pub mod scan; -pub mod token; -pub mod types; -pub mod xata; - -fn cli() -> Command { - Command::new("dropbox") - .version(env!("CARGO_PKG_VERSION")) - .about("Rocksky Dropbox Service") - .subcommand(Command::new("scan").about("Scan Dropbox Music Folder")) - .subcommand(Command::new("serve").about("Serve Rocksky Dropbox API")) -} - -#[tokio::main] -async fn main() -> Result<(), Box> { - dotenv().ok(); - - let args = cli().get_matches(); - - match args.subcommand() { - Some(("scan", _)) => scan().await?, - Some(("serve", _)) => serve().await?, - _ => serve().await?, - } - - Ok(()) -} diff --git a/crates/googledrive/src/main.rs b/crates/googledrive/src/main.rs deleted file mode 100644 index bd6ceb21..00000000 --- a/crates/googledrive/src/main.rs +++ /dev/null @@ -1,37 +0,0 @@ -use clap::Command; -use cmd::{scan::scan, serve::serve}; -use dotenv::dotenv; - -pub mod client; -pub mod cmd; -pub mod consts; -pub mod crypto; -pub mod handlers; -pub mod repo; -pub mod scan; -pub mod token; -pub mod types; -pub mod xata; - -fn cli() -> Command { - Command::new("googledrive") - .version(env!("CARGO_PKG_VERSION")) - .about("Rocksky Google Drive Service") - .subcommand(Command::new("scan").about("Scan Google Drive Music Folder")) - .subcommand(Command::new("serve").about("Serve Rocksky Google Drive API")) -} - -#[tokio::main] -async fn main() -> Result<(), Box> { - dotenv().ok(); - - let args = cli().get_matches(); - - match args.subcommand() { - Some(("scan", _)) => scan().await?, - Some(("serve", _)) => serve().await?, - _ => serve().await?, - } - - Ok(()) -} diff --git a/crates/jetstream/src/main.rs b/crates/jetstream/src/main.rs deleted file mode 100644 index d4c5aee7..00000000 --- a/crates/jetstream/src/main.rs +++ /dev/null @@ -1,37 +0,0 @@ -use std::{env, sync::Arc}; - -use dotenv::dotenv; -use subscriber::ScrobbleSubscriber; -use tokio::sync::Mutex; - -use crate::webhook_worker::AppState; - -pub mod profile; -pub mod repo; -pub mod subscriber; -pub mod types; -pub mod webhook; -pub mod webhook_worker; -pub mod xata; - -#[tokio::main] -async fn main() -> Result<(), anyhow::Error> { - dotenv()?; - let jetstream_server = env::var("JETSTREAM_SERVER") - .unwrap_or_else(|_| "wss://jetstream2.us-west.bsky.network".to_string()); - let url = format!( - "{}/subscribe?wantedCollections=app.rocksky.*", - jetstream_server - ); - let subscriber = ScrobbleSubscriber::new(&url); - - let redis_url = env::var("REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string()); - let redis = redis::Client::open(redis_url)?; - let queue_key = - env::var("WEBHOOK_QUEUE_KEY").unwrap_or_else(|_| "rocksky:webhook_queue".to_string()); - - let state = Arc::new(Mutex::new(AppState { redis, queue_key })); - - subscriber.run(state).await?; - Ok(()) -} diff --git a/crates/playlists/src/main.rs b/crates/playlists/src/main.rs deleted file mode 100644 index 5b845cd7..00000000 --- a/crates/playlists/src/main.rs +++ /dev/null @@ -1,65 +0,0 @@ -use core::{create_tables, find_spotify_users, load_users, save_playlists}; -use std::{ - env, - sync::{Arc, Mutex}, -}; - -use anyhow::Error; -use async_nats::connect; -use dotenv::dotenv; -use duckdb::Connection; -use owo_colors::OwoColorize; -use rocksky_playlists::subscriber::subscribe; -use spotify::get_user_playlists; -use sqlx::postgres::PgPoolOptions; - -pub mod core; -pub mod crypto; -pub mod spotify; -pub mod types; -pub mod xata; - -#[tokio::main] -async fn main() -> Result<(), Error> { - dotenv().ok(); - - let conn = Connection::open("./rocksky-playlists.ddb")?; - let conn = Arc::new(Mutex::new(conn)); - create_tables(conn.clone())?; - - subscribe(conn.clone()).await?; - - let pool = PgPoolOptions::new() - .max_connections(5) - .connect(&env::var("XATA_POSTGRES_URL")?) - .await?; - let users = find_spotify_users(&pool, 0, 100).await?; - - load_users(conn.clone(), &pool).await?; - - sqlx::query(r#" - CREATE UNIQUE INDEX IF NOT EXISTS user_playlists_unique_index ON user_playlists (user_id, playlist_id) - "#) - .execute(&pool) - .await?; - let conn = conn.clone(); - - let addr = env::var("NATS_URL").unwrap_or_else(|_| "nats://localhost:4222".to_string()); - let nc = connect(&addr).await?; - let nc = Arc::new(Mutex::new(nc)); - println!("Connected to NATS server at {}", addr.bright_green()); - - for user in users { - let token = user.1.clone(); - let did = user.2.clone(); - let user_id = user.3.clone(); - let playlists = get_user_playlists(token).await?; - save_playlists(&pool, conn.clone(), nc.clone(), playlists, &user_id, &did).await?; - } - - println!("Done!"); - - loop { - tokio::time::sleep(tokio::time::Duration::from_secs(1)).await; - } -} diff --git a/crates/scrobbler/src/main.rs b/crates/scrobbler/src/main.rs deleted file mode 100644 index 61dd9cd8..00000000 --- a/crates/scrobbler/src/main.rs +++ /dev/null @@ -1,105 +0,0 @@ -pub mod auth; -pub mod cache; -pub mod crypto; -pub mod handlers; -pub mod listenbrainz; -pub mod musicbrainz; -pub mod params; -pub mod repo; -pub mod response; -pub mod rocksky; -pub mod scrobbler; -pub mod signature; -pub mod spotify; -pub mod types; -pub mod xata; - -use actix_limitation::{Limiter, RateLimiter}; -use actix_session::SessionExt as _; -use actix_web::{ - dev::ServiceRequest, - web::{self, Data}, - App, HttpServer, -}; -use anyhow::Error; -use cache::Cache; -use dotenv::dotenv; -use owo_colors::OwoColorize; -use sqlx::postgres::PgPoolOptions; -use std::{env, sync::Arc, time::Duration}; - -pub const BANNER: &str = r#" - ___ ___ _____ __ __ __ - / | __ ______/ (_)___ / ___/______________ / /_ / /_ / /__ _____ - / /| |/ / / / __ / / __ \ \__ \/ ___/ ___/ __ \/ __ \/ __ \/ / _ \/ ___/ - / ___ / /_/ / /_/ / / /_/ / ___/ / /__/ / / /_/ / /_/ / /_/ / / __/ / -/_/ |_\__,_/\__,_/_/\____/ /____/\___/_/ \____/_.___/_.___/_/\___/_/ - - This is the Rocksky Scrobbler API compatible with Last.fm AudioScrobbler API -"#; - -#[tokio::main] -async fn main() -> Result<(), Error> { - dotenv().ok(); - - println!("{}", BANNER.magenta()); - - let cache = Cache::new()?; - - let pool = PgPoolOptions::new() - .max_connections(5) - .connect(&env::var("XATA_POSTGRES_URL")?) - .await?; - let conn = Arc::new(pool); - - let host = env::var("SCROBBLE_HOST").unwrap_or_else(|_| "127.0.0.1".to_string()); - let port = env::var("SCROBBLE_PORT") - .unwrap_or_else(|_| "7882".to_string()) - .parse::() - .unwrap_or(7882); - - tracing::info!( - url = %format!("http://{}:{}", host, port).bright_green(), - "Starting Scrobble server @" - ); - - let limiter = web::Data::new( - Limiter::builder("redis://127.0.0.1") - .key_by(|req: &ServiceRequest| { - req.get_session() - .get(&"session-id") - .unwrap_or_else(|_| req.cookie(&"rate-api-id").map(|c| c.to_string())) - }) - .limit(100) - .period(Duration::from_secs(60)) // 60 minutes - .build() - .unwrap(), - ); - - HttpServer::new(move || { - App::new() - .wrap(RateLimiter::default()) - .app_data(limiter.clone()) - .app_data(Data::new(conn.clone())) - .app_data(Data::new(cache.clone())) - .service(handlers::handle_methods) - .service(handlers::handle_nowplaying) - .service(handlers::handle_submission) - .service(listenbrainz::handlers::handle_submit_listens) - .service(listenbrainz::handlers::handle_validate_token) - .service(listenbrainz::handlers::handle_search_users) - .service(listenbrainz::handlers::handle_get_playing_now) - .service(listenbrainz::handlers::handle_get_listens) - .service(listenbrainz::handlers::handle_get_listen_count) - .service(listenbrainz::handlers::handle_get_artists) - .service(listenbrainz::handlers::handle_get_recordings) - .service(listenbrainz::handlers::handle_get_release_groups) - .service(handlers::index) - .service(handlers::handle_get) - }) - .bind((host, port))? - .run() - .await?; - - Ok(()) -} diff --git a/crates/webscrobbler/src/main.rs b/crates/webscrobbler/src/main.rs deleted file mode 100644 index 7c0f3f14..00000000 --- a/crates/webscrobbler/src/main.rs +++ /dev/null @@ -1,12 +0,0 @@ -use anyhow::Error; -use dotenv::dotenv; -use rocksky_webscrobbler::start_server; - -#[tokio::main] -async fn main() -> Result<(), Error> { - dotenv().ok(); - - start_server().await?; - - Ok(()) -}