diff --git a/Cargo.lock b/Cargo.lock index 4ae0127a..ef19792c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -277,13 +277,13 @@ dependencies = [ "clap", "dotenv", "duckdb", - "futures-util", "owo-colors", "polars", "serde", "serde_json", "sqlx", "tokio", + "tokio-stream", ] [[package]] @@ -4519,6 +4519,7 @@ dependencies = [ "futures-core", "pin-project-lite", "tokio", + "tokio-util", ] [[package]] diff --git a/crates/analytics/Cargo.toml b/crates/analytics/Cargo.toml index eb3650c7..a2df20d9 100644 --- a/crates/analytics/Cargo.toml +++ b/crates/analytics/Cargo.toml @@ -27,4 +27,4 @@ anyhow = "1.0.96" polars = "0.46.0" clap = "4.5.31" actix-web = "4.9.0" -futures-util = "0.3.31" +tokio-stream = { version = "0.1.17", features = ["full"] } diff --git a/crates/analytics/src/cmd/serve.rs b/crates/analytics/src/cmd/serve.rs index 7e35cd0b..c03cd77b 100644 --- a/crates/analytics/src/cmd/serve.rs +++ b/crates/analytics/src/cmd/serve.rs @@ -1,16 +1,21 @@ use std::env; -use actix_web::{get, post, web::{self, Data}, App, HttpRequest, HttpServer, Responder}; +use actix_web::{get, post, web::{self, Data}, App, HttpRequest, HttpResponse, HttpServer, Responder}; use duckdb::Connection; use anyhow::Error; use owo_colors::OwoColorize; +use serde_json::json; use std::sync::{Arc, Mutex}; -use crate::handlers::handle; +use crate::{handlers::handle, subscriber::subscribe}; +// return json response #[get("/")] -async fn index(_req: HttpRequest) -> String { - "Hello world!".to_owned() +async fn index(_req: HttpRequest) -> HttpResponse { + HttpResponse::Ok().json(json!({ + "server": "Rocksky Analytics Server", + "version": "0.1.0", + })) } #[post("/{method}")] @@ -28,6 +33,8 @@ async fn call_method( pub async fn serve(conn: Arc>) -> Result<(), Error> { + subscribe(conn.clone()).await?; + let host = env::var("ANALYTICS_HOST").unwrap_or_else(|_| "127.0.0.1".to_string()); let port = env::var("ANALYTICS_PORT").unwrap_or_else(|_| "7879".to_string()); let addr = format!("{}:{}", host, port); diff --git a/crates/analytics/src/cmd/sync.rs b/crates/analytics/src/cmd/sync.rs index 88569032..05408737 100644 --- a/crates/analytics/src/cmd/sync.rs +++ b/crates/analytics/src/cmd/sync.rs @@ -14,6 +14,7 @@ pub async fn sync(conn: Arc>, pool: &Pool) -> Result load_album_tracks(conn.clone(), pool).await?; load_loved_tracks(conn.clone(), pool).await?; load_artist_tracks(conn.clone(), pool).await?; + load_artist_albums(conn.clone(), pool).await?; load_user_albums(conn.clone(), pool).await?; load_user_artists(conn.clone(), pool).await?; load_user_tracks(conn.clone(), pool).await?; diff --git a/crates/analytics/src/core.rs b/crates/analytics/src/core.rs index 1f1c475e..9870d62d 100644 --- a/crates/analytics/src/core.rs +++ b/crates/analytics/src/core.rs @@ -149,6 +149,14 @@ pub async fn create_tables(conn: &Connection) -> Result<(), Error> { FOREIGN KEY (artist_id) REFERENCES artists(id), FOREIGN KEY (track_id) REFERENCES tracks(id), ); + CREATE TABLE IF NOT EXISTS artist_albums ( + id VARCHAR PRIMARY KEY, + artist_id VARCHAR, + album_id VARCHAR, + created_at TIMESTAMP, + FOREIGN KEY (artist_id) REFERENCES artists(id), + FOREIGN KEY (album_id) REFERENCES albums(id), + ); CREATE TABLE IF NOT EXISTS album_tracks ( id VARCHAR PRIMARY KEY, album_id VARCHAR, @@ -438,13 +446,15 @@ pub async fn load_scrobbles(conn: Arc>, pool: &Pool) artist_id, uri, created_at - ) VALUES (?, + ) VALUES ( ?, ?, ?, ?, ?, - ?)", + ?, + ? + )", params![ scrobble.xata_id, scrobble.user_id, @@ -560,6 +570,35 @@ pub async fn load_artist_tracks(conn: Arc>, pool: &Pool>, pool: &Pool) -> Result<(), Error> { + let conn = conn.lock().unwrap(); + let artist_albums: Vec = sqlx::query_as(r#" + SELECT * FROM artist_albums + "#) + .fetch_all(pool) + .await?; + + for (i, artist_album) in artist_albums.clone().into_iter().enumerate() { + println!("artist_albums {} - {} - {}", i, artist_album.artist_id.bright_green(), artist_album.album_id); + match conn.execute( + "INSERT INTO artist_albums (id, artist_id, album_id, created_at) VALUES (?, ?, ?, ?)", + params![ + artist_album.xata_id, + artist_album.artist_id, + artist_album.album_id, + artist_album.xata_createdat, + ], + ) { + Ok(_) => (), + Err(e) => println!("error: {}", e), + } + } + + println!("artist_albums: {:?}", artist_albums.len()); + Ok(()) +} + pub async fn load_user_albums(conn: Arc>, pool: &Pool) -> Result<(), Error> { let conn = conn.lock().unwrap(); let user_albums: Vec = sqlx::query_as(r#" diff --git a/crates/analytics/src/handlers/albums.rs b/crates/analytics/src/handlers/albums.rs index 571ab49e..57d163ec 100644 --- a/crates/analytics/src/handlers/albums.rs +++ b/crates/analytics/src/handlers/albums.rs @@ -1,10 +1,10 @@ use std::sync::{Arc, Mutex}; use actix_web::{web, HttpRequest, HttpResponse}; -use analytics::types::album::{Album, GetAlbumsParams, GetTopAlbumsParams}; +use analytics::types::{album::{Album, GetAlbumTracksParams, GetAlbumsParams, GetTopAlbumsParams}, track::Track}; use duckdb::Connection; use anyhow::Error; -use futures_util::StreamExt; +use tokio_stream::StreamExt; use crate::read_payload; @@ -23,7 +23,7 @@ pub async fn get_albums(payload: &mut web::Payload, _req: &HttpRequest, conn: Ar SELECT a.* FROM user_albums ua LEFT JOIN albums a ON ua.album_id = a.id LEFT JOIN users u ON ua.user_id = u.id - WHERE u.did = ? + WHERE u.did = ? OR u.handle = ? ORDER BY a.title ASC OFFSET ? LIMIT ?; "#)? }, @@ -34,7 +34,7 @@ pub async fn get_albums(payload: &mut web::Payload, _req: &HttpRequest, conn: Ar match did { Some(did) => { - let albums_iter = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + let albums_iter = stmt.query_map([&did, &did, &limit.to_string(), &offset.to_string()], |row| { Ok(Album { id: row.get(0)?, title: row.get(1)?, @@ -101,7 +101,8 @@ pub async fn get_top_albums(payload: &mut web::Payload, _req: &HttpRequest, conn a.album_art AS album_art, a.release_date, a.year, - a.uri AS uri, + a.uri, + a.sha256, COUNT(*) AS play_count, COUNT(DISTINCT s.user_id) AS unique_listeners FROM @@ -112,9 +113,9 @@ pub async fn get_top_albums(payload: &mut web::Payload, _req: &HttpRequest, conn artists ar ON a.artist_uri = ar.uri LEFT JOIN users u ON s.user_id = u.id - WHERE s.album_id IS NOT NULL AND u.did = ? + WHERE s.album_id IS NOT NULL AND (u.did = ? OR u.handle = ?) GROUP BY - s.album_id, a.title, ar.name, a.release_date, a.year, a.uri, a.album_art + s.album_id, a.title, ar.name, a.release_date, a.year, a.uri, a.album_art, a.sha256 ORDER BY play_count DESC OFFSET ? @@ -128,7 +129,8 @@ pub async fn get_top_albums(payload: &mut web::Payload, _req: &HttpRequest, conn a.album_art AS album_art, a.release_date, a.year, - a.uri AS uri, + a.uri, + a.sha256, COUNT(*) AS play_count, COUNT(DISTINCT s.user_id) AS unique_listeners FROM @@ -138,7 +140,7 @@ pub async fn get_top_albums(payload: &mut web::Payload, _req: &HttpRequest, conn LEFT JOIN artists ar ON a.artist_uri = ar.uri WHERE s.album_id IS NOT NULL GROUP BY - s.album_id, a.title, ar.name, a.release_date, a.year, a.uri, a.album_art + s.album_id, a.title, ar.name, a.release_date, a.year, a.uri, a.album_art, a.sha256 ORDER BY play_count DESC OFFSET ? @@ -148,7 +150,7 @@ pub async fn get_top_albums(payload: &mut web::Payload, _req: &HttpRequest, conn match did { Some(did) => { - let albums = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + let albums = stmt.query_map([&did, &did, &limit.to_string(), &offset.to_string()], |row| { Ok(Album { id: row.get(0)?, title: row.get(1)?, @@ -157,8 +159,9 @@ pub async fn get_top_albums(payload: &mut web::Payload, _req: &HttpRequest, conn release_date: row.get(4)?, year: row.get(5)?, uri: row.get(6)?, - play_count: Some(row.get(7)?), - unique_listeners: Some(row.get(8)?), + sha256: row.get(7)?, + play_count: Some(row.get(8)?), + unique_listeners: Some(row.get(9)?), ..Default::default() }) })?; @@ -175,8 +178,9 @@ pub async fn get_top_albums(payload: &mut web::Payload, _req: &HttpRequest, conn release_date: row.get(4)?, year: row.get(5)?, uri: row.get(6)?, - play_count: Some(row.get(7)?), - unique_listeners: Some(row.get(8)?), + sha256: row.get(7)?, + play_count: Some(row.get(8)?), + unique_listeners: Some(row.get(9)?), ..Default::default() }) })?; @@ -184,4 +188,66 @@ pub async fn get_top_albums(payload: &mut web::Payload, _req: &HttpRequest, conn Ok(HttpResponse::Ok().json(web::Json(albums?))) } } -} \ No newline at end of file +} + +pub async fn get_album_tracks(payload: &mut web::Payload, _req: &HttpRequest, conn: Arc>) -> Result { + let body = read_payload!(payload); + let params = serde_json::from_slice::(&body)?; + let conn = conn.lock().unwrap(); + let mut stmt = conn.prepare(r#" + SELECT + t.id, + t.title, + t.artist, + t.album_artist, + t.album, + t.uri, + t.album_art, + t.duration, + t.disc_number, + t.track_number, + t.artist_uri, + t.album_uri, + t.sha256, + t.copyright_message, + t.label, + t.created_at, + COUNT(*) AS play_count, + COUNT(DISTINCT s.user_id) AS unique_listeners + FROM album_tracks at + LEFT JOIN tracks t ON at.track_id = t.id + LEFT JOIN albums a ON at.album_id = a.id + LEFT JOIN scrobbles s ON s.track_id = t.id + WHERE at.album_id = ? OR a.uri = ? + GROUP BY + t.id, t.title, t.artist, t.album_artist, t.album, t.uri, t.album_art, t.duration, t.disc_number, t.track_number, t.artist_uri, t.album_uri, t.sha256, t.copyright_message, t.label, t.created_at + ORDER BY t.track_number ASC; + "#)?; + + let tracks = stmt.query_map([¶ms.album_id, ¶ms.album_id], |row| { + Ok(Track { + id: row.get(0)?, + title: row.get(1)?, + artist: row.get(2)?, + album_artist: row.get(3)?, + album: row.get(4)?, + uri: row.get(5)?, + album_art: row.get(6)?, + duration: row.get(7)?, + disc_number: row.get(8)?, + track_number: row.get(9)?, + artist_uri: row.get(10)?, + album_uri: row.get(11)?, + sha256: row.get(12)?, + copyright_message: row.get(13)?, + label: row.get(14)?, + created_at: row.get(15)?, + play_count: Some(row.get(16)?), + unique_listeners: Some(row.get(17)?), + ..Default::default() + }) + })?; + + let tracks: Result, _> = tracks.collect(); + Ok(HttpResponse::Ok().json(web::Json(tracks?))) +} diff --git a/crates/analytics/src/handlers/artists.rs b/crates/analytics/src/handlers/artists.rs index db0cfa6c..e3875a74 100644 --- a/crates/analytics/src/handlers/artists.rs +++ b/crates/analytics/src/handlers/artists.rs @@ -1,12 +1,12 @@ use std::sync::{Arc, Mutex}; use actix_web::{web, HttpRequest, HttpResponse}; -use analytics::types::artist::{Artist, GetTopArtistsParams}; +use analytics::types::{album::Album, artist::{Artist, GetArtistAlbumsParams, GetArtistTracksParams, GetArtistsParams, GetTopArtistsParams}, track::Track}; use duckdb::Connection; use anyhow::Error; -use futures_util::StreamExt; +use tokio_stream::StreamExt; -use crate::{read_payload, types::artist::GetArtistsParams}; +use crate::read_payload; pub async fn get_artists(payload: &mut web::Payload, _req: &HttpRequest, conn: Arc>) -> Result { let body = read_payload!(payload); @@ -23,7 +23,7 @@ pub async fn get_artists(payload: &mut web::Payload, _req: &HttpRequest, conn: A SELECT a.* FROM user_artists ua LEFT JOIN artists a ON ua.artist_id = a.id LEFT JOIN users u ON ua.user_id = u.id - WHERE u.did = ? + WHERE u.did = ? OR u.handle = ? ORDER BY a.name ASC OFFSET ? LIMIT ?; "#)? }, @@ -34,7 +34,7 @@ pub async fn get_artists(payload: &mut web::Payload, _req: &HttpRequest, conn: A match did { Some(did) => { - let artists = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + let artists = stmt.query_map([&did, &did, &limit.to_string(), &offset.to_string()], |row| { Ok(Artist { id: row.get(0)?, name: row.get(1)?, @@ -111,7 +111,7 @@ pub async fn get_top_artists(payload: &mut web::Payload, _req: &HttpRequest, con LEFT JOIN users u ON s.user_id = u.id WHERE - s.artist_id IS NOT NULL AND u.did = ? + s.artist_id IS NOT NULL AND (u.did = ? OR u.handle = ?) GROUP BY s.artist_id, ar.name, ar.uri, ar.picture, ar.sha256 ORDER BY @@ -148,7 +148,7 @@ pub async fn get_top_artists(payload: &mut web::Payload, _req: &HttpRequest, con match did { Some(did) => { - let artists = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + let artists = stmt.query_map([&did, &did, &limit.to_string(), &offset.to_string()], |row| { Ok(Artist { id: row.get(0)?, name: row.get(1)?, @@ -196,4 +196,124 @@ pub async fn get_top_artists(payload: &mut web::Payload, _req: &HttpRequest, con Ok(HttpResponse::Ok().json(artists?)) } } -} \ No newline at end of file +} + +pub async fn get_artist_tracks(payload: &mut web::Payload, _req: &HttpRequest, conn: Arc>) -> Result { + let body = read_payload!(payload); + let params = serde_json::from_slice::(&body)?; + let pagination = params.pagination.unwrap_or_default(); + let offset = pagination.skip.unwrap_or(0); + let limit = pagination.take.unwrap_or(20); + let conn = conn.lock().unwrap(); + + let mut stmt = conn.prepare(r#" + SELECT + t.id, + t.title, + t.artist, + t.album_artist, + t.album, + t.uri, + t.album_art, + t.duration, + t.disc_number, + t.track_number, + t.artist_uri, + t.album_uri, + t.sha256, + t.copyright_message, + t.label, + t.created_at, + COUNT(*) AS play_count, + COUNT(DISTINCT s.user_id) AS unique_listeners + FROM artist_tracks at + LEFT JOIN tracks t ON at.track_id = t.id + LEFT JOIN artists a ON at.artist_id = a.id + LEFT JOIN scrobbles s ON s.track_id = t.id + WHERE at.artist_id = ? OR a.uri = ? + GROUP BY + t.id, t.title, t.artist, t.album_artist, t.album, t.uri, t.album_art, t.duration, t.disc_number, t.track_number, t.artist_uri, t.album_uri, t.sha256, t.copyright_message, t.label, t.created_at + ORDER BY play_count DESC + OFFSET ? + LIMIT ?; + "#)?; + + let tracks = stmt.query_map([¶ms.artist_id, ¶ms.artist_id, &limit.to_string(), &offset.to_string()], |row| { + Ok(Track { + id: row.get(0)?, + title: row.get(1)?, + artist: row.get(2)?, + album_artist: row.get(3)?, + album: row.get(4)?, + uri: row.get(5)?, + album_art: row.get(6)?, + duration: row.get(7)?, + disc_number: row.get(8)?, + track_number: row.get(9)?, + artist_uri: row.get(10)?, + album_uri: row.get(11)?, + sha256: row.get(12)?, + copyright_message: row.get(13)?, + label: row.get(14)?, + created_at: row.get(15)?, + play_count: Some(row.get(16)?), + unique_listeners: Some(row.get(17)?), + ..Default::default() + }) + })?; + + let tracks: Result, _> = tracks.collect(); + Ok(HttpResponse::Ok().json(tracks?)) +} + +pub async fn get_artist_albums(payload: &mut web::Payload, _req: &HttpRequest, conn: Arc>) -> Result { + let body = read_payload!(payload); + let params = serde_json::from_slice::(&body)?; + let conn = conn.lock().unwrap(); + + let mut stmt = conn.prepare(r#" + SELECT + al.id, + al.title, + al.artist, + al.album_art, + al.release_date, + al.year, + al.uri, + al.sha256, + al.artist_uri, + COUNT(*) AS play_count, + COUNT(DISTINCT s.user_id) AS unique_listeners + FROM + artist_albums aa + LEFT JOIN artists ar ON aa.artist_id = ar.id + LEFT JOIN albums al ON aa.album_id = al.id + LEFT JOIN scrobbles s ON aa.album_id = s.album_id + WHERE ar.id = ? OR ar.uri = ? + GROUP BY al.id, al.title, al.artist, al.album_art, al.release_date, al.year, al.uri, al.sha256, al.artist_uri + ORDER BY play_count DESC; + "#)?; + + let albums = stmt.query_map([¶ms.artist_id, ¶ms.artist_id], |row| { + Ok(Album { + id: row.get(0)?, + title: row.get(1)?, + artist: row.get(2)?, + release_date: row.get(3)?, + album_art: row.get(4)?, + year: row.get(5)?, + spotify_link: None, + tidal_link: None, + youtube_link: None, + apple_music_link: None, + sha256: row.get(7)?, + uri: row.get(6)?, + artist_uri: row.get(8)?, + play_count: Some(row.get(9)?), + unique_listeners: Some(row.get(10)?), + }) + })?; + + let albums: Result, _> = albums.collect(); + Ok(HttpResponse::Ok().json(albums?)) +} diff --git a/crates/analytics/src/handlers/mod.rs b/crates/analytics/src/handlers/mod.rs index f104148e..3764aa1e 100644 --- a/crates/analytics/src/handlers/mod.rs +++ b/crates/analytics/src/handlers/mod.rs @@ -1,11 +1,11 @@ use std::sync::{Arc, Mutex}; use actix_web::{web, HttpRequest, HttpResponse}; -use albums::{get_albums, get_top_albums}; -use artists::{get_artists, get_top_artists}; +use albums::{get_album_tracks, get_albums, get_top_albums}; +use artists::{get_artist_albums, get_artist_tracks, get_artists, get_top_artists}; use duckdb::Connection; use scrobbles::get_scrobbles; -use stats::{get_scrobbles_per_day, get_scrobbles_per_month, get_scrobbles_per_year, get_stats}; +use stats::{get_album_scrobbles, get_artist_scrobbles, get_scrobbles_per_day, get_scrobbles_per_month, get_scrobbles_per_year, get_stats, get_track_scrobbles}; use tracks::{get_loved_tracks, get_top_tracks, get_tracks}; use anyhow::Error; @@ -30,7 +30,6 @@ macro_rules! read_payload { }}; } - pub async fn handle(method: &str, payload: &mut web::Payload, req: &HttpRequest, conn: Arc>) -> Result { match method { "library.getAlbums" => get_albums(payload, req, conn.clone()).await, @@ -45,6 +44,12 @@ pub async fn handle(method: &str, payload: &mut web::Payload, req: &HttpRequest, "library.getScrobblesPerDay" => get_scrobbles_per_day(payload, req, conn.clone()).await, "library.getScrobblesPerMonth" => get_scrobbles_per_month(payload, req, conn.clone()).await, "library.getScrobblesPerYear" => get_scrobbles_per_year(payload, req, conn.clone()).await, + "library.getAlbumScrobbles" => get_album_scrobbles(payload, req, conn.clone()).await, + "library.getArtistScrobbles" => get_artist_scrobbles(payload, req, conn.clone()).await, + "library.getTrackScrobbles" => get_track_scrobbles(payload, req, conn.clone()).await, + "library.getAlbumTracks" => get_album_tracks(payload, req, conn.clone()).await, + "library.getArtistAlbums" => get_artist_albums(payload, req, conn.clone()).await, + "library.getArtistTracks" => get_artist_tracks(payload, req, conn.clone()).await, _ => return Err(anyhow::anyhow!("Method not found")), } } diff --git a/crates/analytics/src/handlers/scrobbles.rs b/crates/analytics/src/handlers/scrobbles.rs index f12542df..f041f9b2 100644 --- a/crates/analytics/src/handlers/scrobbles.rs +++ b/crates/analytics/src/handlers/scrobbles.rs @@ -4,7 +4,7 @@ use actix_web::{web, HttpRequest, HttpResponse}; use analytics::types::scrobble::{GetScrobblesParams, ScrobbleTrack}; use duckdb::Connection; use anyhow::Error; -use futures_util::StreamExt; +use tokio_stream::StreamExt; use crate::read_payload; @@ -38,7 +38,7 @@ pub async fn get_scrobbles(payload: &mut web::Payload, _req: &HttpRequest, conn: LEFT JOIN albums al ON s.album_id = al.id LEFT JOIN tracks t ON s.track_id = t.id LEFT JOIN users u ON s.user_id = u.id - WHERE u.did = ? + WHERE u.did = ? OR u.handle = ? GROUP BY s.id, s.created_at, t.id, t.title, t.artist, t.album_artist, t.album, t.album_art, s.uri, t.uri, u.handle, a.uri, al.uri, s.created_at ORDER BY s.created_at DESC OFFSET ? @@ -72,7 +72,7 @@ pub async fn get_scrobbles(payload: &mut web::Payload, _req: &HttpRequest, conn: }; match did { Some(did) => { - let scrobbles = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + let scrobbles = stmt.query_map([&did, &did, &limit.to_string(), &offset.to_string()], |row| { Ok(ScrobbleTrack { id: row.get(0)?, track_id: row.get(1)?, diff --git a/crates/analytics/src/handlers/stats.rs b/crates/analytics/src/handlers/stats.rs index 76bc6a51..c717d881 100644 --- a/crates/analytics/src/handlers/stats.rs +++ b/crates/analytics/src/handlers/stats.rs @@ -1,11 +1,11 @@ use std::sync::{Arc, Mutex}; use actix_web::{web, HttpRequest, HttpResponse}; -use analytics::types::{scrobble::{ScrobblesPerDay, ScrobblesPerMonth, ScrobblesPerYear}, stats::{GetScrobblesPerDayParams, GetScrobblesPerMonthParams, GetScrobblesPerYearParams, GetStatsParams}}; +use analytics::types::{scrobble::{ScrobblesPerDay, ScrobblesPerMonth, ScrobblesPerYear}, stats::{GetAlbumScrobblesParams, GetArtistScrobblesParams, GetScrobblesPerDayParams, GetScrobblesPerMonthParams, GetScrobblesPerYearParams, GetStatsParams, GetTrackScrobblesParams}}; use duckdb::Connection; use anyhow::Error; use serde_json::json; -use futures_util::StreamExt; +use tokio_stream::StreamExt; use crate::read_payload; pub async fn get_stats(payload: &mut web::Payload, _req: &HttpRequest, conn: Arc>) -> Result { @@ -13,20 +13,20 @@ pub async fn get_stats(payload: &mut web::Payload, _req: &HttpRequest, conn: Arc let params = serde_json::from_slice::(&body)?; let conn = conn.lock().unwrap(); - let mut stmt = conn.prepare("SELECT COUNT(*) FROM scrobbles s LEFT JOIN users u ON s.user_id = u.id WHERE u.did = ?")?; - let scrobbles: i64 = stmt.query_row([¶ms.user_did], |row| row.get(0))?; + let mut stmt = conn.prepare("SELECT COUNT(*) FROM scrobbles s LEFT JOIN users u ON s.user_id = u.id WHERE u.did = ? OR u.handle = ?")?; + let scrobbles: i64 = stmt.query_row([¶ms.user_did, ¶ms.user_did], |row| row.get(0))?; - let mut stmt = conn.prepare("SELECT COUNT(*) FROM user_artists LEFT JOIN users u ON user_artists.user_id = u.id WHERE u.did = ?")?; - let artists: i64 = stmt.query_row([¶ms.user_did], |row| row.get(0))?; + let mut stmt = conn.prepare("SELECT COUNT(*) FROM user_artists LEFT JOIN users u ON user_artists.user_id = u.id WHERE u.did = ? OR u.handle = ?")?; + let artists: i64 = stmt.query_row([¶ms.user_did, ¶ms.user_did], |row| row.get(0))?; - let mut stmt = conn.prepare("SELECT COUNT(*) FROM loved_tracks LEFT JOIN users u ON loved_tracks.user_id = u.id WHERE u.did = ?")?; - let loved_tracks: i64 = stmt.query_row([¶ms.user_did], |row| row.get(0))?; + let mut stmt = conn.prepare("SELECT COUNT(*) FROM loved_tracks LEFT JOIN users u ON loved_tracks.user_id = u.id WHERE u.did = ? OR u.handle = ?")?; + let loved_tracks: i64 = stmt.query_row([¶ms.user_did, ¶ms.user_did], |row| row.get(0))?; - let mut stmt = conn.prepare("SELECT COUNT(*) FROM user_albums LEFT JOIN users u ON user_albums.user_id = u.id WHERE u.did = ?")?; - let albums: i64 = stmt.query_row([¶ms.user_did], |row| row.get(0))?; + let mut stmt = conn.prepare("SELECT COUNT(*) FROM user_albums LEFT JOIN users u ON user_albums.user_id = u.id WHERE u.did = ? OR u.handle = ?")?; + let albums: i64 = stmt.query_row([¶ms.user_did, ¶ms.user_did], |row| row.get(0))?; - let mut stmt = conn.prepare("SELECT COUNT(*) FROM user_tracks LEFT JOIN users u ON user_tracks.user_id = u.id WHERE u.did = ?")?; - let tracks: i64 = stmt.query_row([¶ms.user_did], |row| row.get(0))?; + let mut stmt = conn.prepare("SELECT COUNT(*) FROM user_tracks LEFT JOIN users u ON user_tracks.user_id = u.id WHERE u.did = ? OR u.handle = ?")?; + let tracks: i64 = stmt.query_row([¶ms.user_did, ¶ms.user_did], |row| row.get(0))?; Ok(HttpResponse::Ok().json(json!({ "scrobbles": scrobbles, @@ -55,14 +55,14 @@ pub async fn get_scrobbles_per_day(payload: &mut web::Payload, _req: &HttpReques scrobbles LEFT JOIN users u ON scrobbles.user_id = u.id WHERE - u.did = ? + u.did = ? OR u.handle = ? AND created_at BETWEEN ? AND ? GROUP BY date_trunc('day', created_at) ORDER BY date; "#)?; - let scrobbles = stmt.query_map([did, start, end], |row| { + let scrobbles = stmt.query_map([&did, &did, &start, &end], |row| { Ok(ScrobblesPerDay { date: row.get(0)?, count: row.get(1)?, @@ -116,7 +116,7 @@ pub async fn get_scrobbles_per_month(payload: &mut web::Payload, _req: &HttpRequ scrobbles LEFT JOIN users u ON scrobbles.user_id = u.id WHERE - u.did = ? + u.did = ? OR u.handle = ? AND created_at BETWEEN ? AND ? GROUP BY EXTRACT(YEAR FROM created_at), @@ -124,7 +124,7 @@ pub async fn get_scrobbles_per_month(payload: &mut web::Payload, _req: &HttpRequ ORDER BY year_month; "#)?; - let scrobbles = stmt.query_map([did, start, end], |row| { + let scrobbles = stmt.query_map([&did, &did, &start, &end], |row| { Ok(ScrobblesPerMonth { year_month: row.get(0)?, count: row.get(1)?, @@ -179,14 +179,14 @@ pub async fn get_scrobbles_per_year(payload: &mut web::Payload, _req: &HttpReque scrobbles LEFT JOIN users u ON scrobbles.user_id = u.id WHERE - u.did = ? + u.did = ? OR u.handle = ? AND created_at BETWEEN ? AND ? GROUP BY EXTRACT(YEAR FROM created_at) ORDER BY year; "#)?; - let scrobbles = stmt.query_map([did, start, end], |row| { + let scrobbles = stmt.query_map([&did, &did, &start, &end], |row| { Ok(ScrobblesPerYear { year: row.get(0)?, count: row.get(1)?, @@ -220,3 +220,117 @@ pub async fn get_scrobbles_per_year(payload: &mut web::Payload, _req: &HttpReque } } } + +pub async fn get_album_scrobbles(payload: &mut web::Payload, _req: &HttpRequest, conn: Arc>) -> Result { + let body = read_payload!(payload); + let params = serde_json::from_slice::(&body)?; + let start = params.start.unwrap_or(GetAlbumScrobblesParams::default().start.unwrap()); + let end = params.end.unwrap_or(GetAlbumScrobblesParams::default().end.unwrap()); + let conn = conn.lock().unwrap(); + let mut stmt = conn.prepare(r#" + SELECT + date_trunc('day', s.created_at) AS date, + COUNT(s.album_id) AS count + FROM + scrobbles s + LEFT JOIN albums a ON s.album_id = a.id + WHERE + a.id = ? OR a.uri = ? + AND s.created_at BETWEEN ? AND ? + GROUP BY + date_trunc('day', s.created_at) + ORDER BY + date; + "#)?; + let scrobbles = stmt.query_map([ + ¶ms.album_id, + ¶ms.album_id, + &start, + &end + ], |row| { + Ok(ScrobblesPerDay { + date: row.get(0)?, + count: row.get(1)?, + }) + })?; + let scrobbles: Result, _> = scrobbles.collect(); + Ok(HttpResponse::Ok().json(scrobbles?)) +} + +pub async fn get_artist_scrobbles(payload: &mut web::Payload, _req: &HttpRequest, conn: Arc>) -> Result { + let body = read_payload!(payload); + let params = serde_json::from_slice::(&body)?; + let start = params.start.unwrap_or(GetArtistScrobblesParams::default().start.unwrap()); + let end = params.end.unwrap_or(GetArtistScrobblesParams::default().end.unwrap()); + let conn = conn.lock().unwrap(); + + let mut stmt = conn.prepare(r#" + SELECT + date_trunc('day', s.created_at) AS date, + COUNT(s.artist_id) AS count + FROM + scrobbles s + LEFT JOIN artists a ON s.artist_id = a.id + WHERE + a.id = ? OR a.uri = ? + AND s.created_at BETWEEN ? AND ? + GROUP BY + date_trunc('day', s.created_at) + ORDER BY + date; + "#)?; + + let scrobbles = stmt.query_map([ + ¶ms.artist_id, + ¶ms.artist_id, + &start, + &end + ], |row| { + Ok(ScrobblesPerDay { + date: row.get(0)?, + count: row.get(1)?, + }) + })?; + + let scrobbles: Result, _> = scrobbles.collect(); + Ok(HttpResponse::Ok().json(scrobbles?)) +} + +pub async fn get_track_scrobbles(payload: &mut web::Payload, _req: &HttpRequest, conn: Arc>) -> Result { + let body = read_payload!(payload); + let params = serde_json::from_slice::(&body)?; + let start = params.start.unwrap_or(GetTrackScrobblesParams::default().start.unwrap()); + let end = params.end.unwrap_or(GetTrackScrobblesParams::default().end.unwrap()); + let conn = conn.lock().unwrap(); + + let mut stmt = conn.prepare(r#" + SELECT + date_trunc('day', s.created_at) AS date, + COUNT(s.track_id) AS count + FROM + scrobbles s + LEFT JOIN tracks t ON s.track_id = t.id + WHERE + t.id = ? OR t.uri = ? + AND s.created_at BETWEEN ? AND ? + GROUP BY + date_trunc('day', s.created_at) + ORDER BY + date; + "#)?; + + let scrobbles = stmt.query_map([ + ¶ms.track_id, + ¶ms.track_id, + &start, + &end + ], |row| { + Ok(ScrobblesPerDay { + date: row.get(0)?, + count: row.get(1)?, + }) + })?; + + let scrobbles: Result, _> = scrobbles.collect(); + Ok(HttpResponse::Ok().json(scrobbles?)) +} diff --git a/crates/analytics/src/handlers/tracks.rs b/crates/analytics/src/handlers/tracks.rs index 2e283525..7212a974 100644 --- a/crates/analytics/src/handlers/tracks.rs +++ b/crates/analytics/src/handlers/tracks.rs @@ -4,7 +4,7 @@ use actix_web::{web, HttpRequest, HttpResponse}; use analytics::types::track::{GetLovedTracksParams, GetTopTracksParams, GetTracksParams, Track}; use duckdb::Connection; use anyhow::Error; -use futures_util::StreamExt; +use tokio_stream::StreamExt; use crate::read_payload; @@ -48,12 +48,12 @@ pub async fn get_tracks(payload: &mut web::Payload, _req: &HttpRequest, conn: Ar FROM tracks t LEFT JOIN user_tracks ut ON t.id = ut.track_id LEFT JOIN users u ON ut.user_id = u.id - WHERE u.did = ? + WHERE u.did = ? OR u.handle = ? ORDER BY t.title ASC OFFSET ? LIMIT ?; "#)?; - let tracks = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + let tracks = stmt.query_map([&did, &did, &limit.to_string(), &offset.to_string()], |row| { Ok(Track { id: row.get(0)?, title: row.get(1)?, @@ -186,12 +186,12 @@ pub async fn get_loved_tracks(payload: &mut web::Payload, _req: &HttpRequest, co FROM loved_tracks l LEFT JOIN users u ON l.user_id = u.id LEFT JOIN tracks t ON l.track_id = t.id - WHERE u.did = ? + WHERE u.did = ? OR u.handle = ? ORDER BY l.created_at DESC OFFSET ? LIMIT ?; "#)?; - let loved_tracks = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + let loved_tracks = stmt.query_map([&did, &did, &limit.to_string(), &offset.to_string()], |row| { Ok(Track { id: row.get(0)?, title: row.get(1)?, @@ -257,13 +257,13 @@ pub async fn get_top_tracks(payload: &mut web::Payload, _req: &HttpRequest, conn LEFT JOIN artists ar ON s.artist_id = ar.id LEFT JOIN albums a ON s.album_id = a.id LEFT JOIN users u ON s.user_id = u.id - WHERE u.did = ? + WHERE u.did = ? OR u.handle = ? GROUP BY t.id, s.track_id, t.title, ar.name, a.title, t.artist, t.uri, t.album_art, t.duration, t.disc_number, t.track_number, t.artist_uri, t.album_uri, t.created_at, t.sha256, t.album_artist, t.album ORDER BY play_count DESC OFFSET ? LIMIT ?; "#)?; - let top_tracks = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + let top_tracks = stmt.query_map([&did, &did, &limit.to_string(), &offset.to_string()], |row| { Ok(Track { id: row.get(0)?, title: row.get(1)?, diff --git a/crates/analytics/src/main.rs b/crates/analytics/src/main.rs index 870b5348..b8744652 100644 --- a/crates/analytics/src/main.rs +++ b/crates/analytics/src/main.rs @@ -12,6 +12,7 @@ pub mod xata; pub mod cmd; pub mod core; pub mod handlers; +pub mod subscriber; fn cli() -> Command { Command::new("analytics") diff --git a/crates/analytics/src/subscriber/mod.rs b/crates/analytics/src/subscriber/mod.rs new file mode 100644 index 00000000..80757181 --- /dev/null +++ b/crates/analytics/src/subscriber/mod.rs @@ -0,0 +1,481 @@ +use std::{env, sync::{Arc, Mutex}, thread}; +use anyhow::Error; +use async_nats::{connect, Client}; +use duckdb::{params, Connection}; +use owo_colors::OwoColorize; +use tokio_stream::StreamExt; +use types::{LikePayload, ScrobblePayload, UnlikePayload}; + +pub mod types; + +pub async fn subscribe(conn: Arc>) -> Result<(), Error> { + let addr = env::var("NATS_URL").unwrap_or_else(|_| "nats://localhost:4222".to_string()); + let conn = conn.clone(); + let nc = connect(&addr).await?; + println!("Connected to NATS server at {}", addr.bright_green()); + + let nc = Arc::new(Mutex::new(nc)); + on_scrobble(nc.clone(), conn.clone()); + on_like(nc.clone(), conn.clone()); + on_unlike(nc.clone(), conn.clone()); + + Ok(()) +} + +pub fn on_scrobble(nc: Arc>, conn: Arc>) { + thread::spawn(move || { + let rt = tokio::runtime::Runtime::new().unwrap(); + let conn = conn.clone(); + let nc = nc.clone(); + rt.block_on(async { + let nc = nc.lock().unwrap(); + let mut sub = nc.subscribe("rocksky.scrobble".to_string()).await?; + drop(nc); + + while let Some(msg) = sub.next().await { + let data = String::from_utf8(msg.payload.to_vec()).unwrap(); + match serde_json::from_str::(&data) { + Ok(payload) => { + match save_scrobble(conn.clone(), payload.clone()).await { + Ok(_) => println!("Scrobble saved successfully for {}", payload.scrobble.uri.cyan()), + Err(e) => eprintln!("Error saving scrobble: {}", e), + } + }, + Err(e) => { + eprintln!("Error parsing payload: {}", e); + } + } + } + + Ok::<(), Error>(()) + })?; + + Ok::<(), Error>(()) + }); +} + + +pub fn on_like(nc: Arc>, conn: Arc>) { + thread::spawn(move || { + let rt = tokio::runtime::Runtime::new().unwrap(); + let conn = conn.clone(); + let nc = nc.clone(); + rt.block_on(async { + let nc = nc.lock().unwrap(); + let mut sub = nc.subscribe("rocksky.like".to_string()).await?; + drop(nc); + + while let Some(msg) = sub.next().await { + let data = String::from_utf8(msg.payload.to_vec()).unwrap(); + match serde_json::from_str::(&data) { + Ok(payload) => { + match like(conn.clone(), payload.clone()).await { + Ok(_) => println!("Like saved successfully for {}", payload.track_id.xata_id.cyan()), + Err(e) => eprintln!("Error saving like: {}", e), + } + }, + Err(e) => { + eprintln!("Error parsing payload: {}", e); + } + } + } + + Ok::<(), Error>(()) + })?; + + Ok::<(), Error>(()) + }); +} + +pub fn on_unlike(nc: Arc>, conn: Arc>) { + thread::spawn(move || { + let rt = tokio::runtime::Runtime::new().unwrap(); + let conn = conn.clone(); + let nc = nc.clone(); + rt.block_on(async { + let nc = nc.lock().unwrap(); + let mut sub = nc.subscribe("rocksky.unlike".to_string()).await?; + drop(nc); + + while let Some(msg) = sub.next().await { + let data = String::from_utf8(msg.payload.to_vec()).unwrap(); + match serde_json::from_str::(&data) { + Ok(payload) => { + match unlike(conn.clone(), payload.clone()).await { + Ok(_) => println!("Unlike saved successfully for {}", payload.track_id.xata_id.cyan()), + Err(e) => eprintln!("Error saving unlike: {}", e), + } + }, + Err(e) => { + eprintln!("Error parsing payload: {}", e); + } + } + } + + Ok::<(), Error>(()) + })?; + + Ok::<(), Error>(()) + }); +} + +pub async fn save_scrobble(conn: Arc>, payload: ScrobblePayload) -> Result<(), Error> { + let conn = conn.lock().unwrap(); + + match conn.execute( + "INSERT INTO artists ( + id, + name, + biography, + born, + born_in, + died, + picture, + sha256, + spotify_link, + tidal_link, + youtube_link, + apple_music_link, + uri + ) VALUES ( + ?, + ?, + ?, + ?, + ?, + ?, + ?, + ?, + ?, + ?, + ?, + ?, + ? + )", + params![ + payload.scrobble.artist_id.xata_id, + payload.scrobble.artist_id.name, + payload.scrobble.artist_id.biography, + payload.scrobble.artist_id.born, + payload.scrobble.artist_id.born_in, + payload.scrobble.artist_id.died, + payload.scrobble.artist_id.picture, + payload.scrobble.artist_id.sha256, + payload.scrobble.artist_id.spotify_link, + payload.scrobble.artist_id.tidal_link, + payload.scrobble.artist_id.youtube_link, + payload.scrobble.artist_id.apple_music_link, + payload.scrobble.artist_id.uri, + ], + ) { + Ok(_) => (), + Err(e) => { + if !e.to_string().contains("violates primary key constraint") { + println!("[artists] error: {}", e); + return Err(e.into()); + } + } + } + + match conn.execute( + "INSERT INTO albums ( + id, + title, + artist, + release_date, + album_art, + year, + spotify_link, + tidal_link, + youtube_link, + apple_music_link, + sha256, + uri, + artist_uri + ) VALUES ( + ?, + ?, + ?, + ?, + ?, + ?, + ?, + ?, + ?, + ?, + ?, + ?, + ? + )", + params![ + payload.scrobble.album_id.xata_id, + payload.scrobble.album_id.title, + payload.scrobble.album_id.artist, + payload.scrobble.album_id.release_date, + payload.scrobble.album_id.album_art, + payload.scrobble.album_id.year, + payload.scrobble.album_id.spotify_link, + payload.scrobble.album_id.tidal_link, + payload.scrobble.album_id.youtube_link, + payload.scrobble.album_id.apple_music_link, + payload.scrobble.album_id.sha256, + payload.scrobble.album_id.uri, + payload.scrobble.album_id.artist_uri, + ], + ) { + Ok(_) => (), + Err(e) => { + if !e.to_string().contains("violates primary key constraint") { + println!("[albums] error: {}", e); + return Err(e.into()); + } + }, + } + + match conn.execute( + "INSERT INTO tracks ( + id, + title, + artist, + album_artist, + album_art, + album, + track_number, + duration, + mb_id, + youtube_link, + spotify_link, + tidal_link, + apple_music_link, + sha256, + lyrics, + composer, + genre, + disc_number, + copyright_message, + label, + uri, + artist_uri, + album_uri, + created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + params![ + payload.scrobble.track_id.xata_id, + payload.scrobble.track_id.title, + payload.scrobble.track_id.artist, + payload.scrobble.track_id.album_artist, + payload.scrobble.track_id.album_art, + payload.scrobble.track_id.album, + payload.scrobble.track_id.track_number, + payload.scrobble.track_id.duration, + payload.scrobble.track_id.mb_id, + payload.scrobble.track_id.youtube_link, + payload.scrobble.track_id.spotify_link, + payload.scrobble.track_id.tidal_link, + payload.scrobble.track_id.apple_music_link, + payload.scrobble.track_id.sha256, + payload.scrobble.track_id.lyrics, + payload.scrobble.track_id.composer, + payload.scrobble.track_id.genre, + payload.scrobble.track_id.disc_number, + payload.scrobble.track_id.copyright_message, + payload.scrobble.track_id.label, + payload.scrobble.track_id.uri, + payload.scrobble.track_id.artist_uri, + payload.scrobble.track_id.album_uri, + payload.scrobble.track_id.xata_createdat, + ], + ) { + Ok(_) => (), + Err(e) => { + if !e.to_string().contains("violates primary key constraint") { + println!("[tracks] error: {}", e); + return Err(e.into()); + } + } + } + + match conn.execute( + "INSERT INTO album_tracks ( + id, + album_id, + track_id + ) VALUES (?, + ?, + ?)", + params![ + payload.album_track.xata_id, + payload.album_track.album_id.xata_id, + payload.album_track.track_id.xata_id, + ], + ) { + Ok(_) => (), + Err(e) => { + if !e.to_string().contains("violates primary key constraint") { + println!("[album_tracks] error: {}", e); + return Err(e.into()); + } + } + } + + match conn.execute( + "INSERT INTO artist_tracks (id, artist_id, track_id, created_at) VALUES (?, ?, ?, ?)", + params![ + payload.artist_track.xata_id, + payload.artist_track.artist_id.xata_id, + payload.artist_track.track_id.xata_id, + payload.artist_track.xata_createdat, + ], + ) { + Ok(_) => (), + Err(e) => { + if !e.to_string().contains("violates primary key constraint") { + println!("[artist_tracks] error: {}", e); + return Err(e.into()); + } + } + } + + match conn.execute( + "INSERT INTO user_albums (id, user_id, album_id, created_at) VALUES (?, ?, ?, ?)", + params![ + payload.user_album.xata_id, + payload.user_album.user_id.xata_id, + payload.user_album.album_id.xata_id, + payload.user_album.xata_createdat, + ], + ) { + Ok(_) => (), + Err(e) => { + if !e.to_string().contains("violates primary key constraint") { + println!("[user_albums] error: {}", e); + return Err(e.into()); + } + } + } + + match conn.execute( + "INSERT INTO user_artists (id, user_id, artist_id, created_at) VALUES (?, ?, ?, ?)", + params![ + payload.user_artist.xata_id, + payload.user_artist.user_id.xata_id, + payload.user_artist.artist_id.xata_id, + payload.user_artist.xata_createdat, + ], + ) { + Ok(_) => (), + Err(e) => { + if !e.to_string().contains("violates primary key constraint") { + println!("[user_artists] error: {}", e); + return Err(e.into()); + } + } + } + + match conn.execute( + "INSERT INTO user_tracks (id, user_id, track_id, created_at) VALUES (?, ?, ?, ?)", + params![ + payload.user_track.xata_id, + payload.user_track.user_id.xata_id, + payload.user_track.track_id.xata_id, + payload.user_track.xata_createdat, + ], + ) { + Ok(_) => (), + Err(e) => { + if !e.to_string().contains("violates primary key constraint") { + println!("[user_tracks] error: {}", e); + return Err(e.into()); + } + } + } + + match conn.execute( + "INSERT INTO scrobbles ( + id, + user_id, + track_id, + album_id, + artist_id, + uri, + created_at + ) VALUES ( + ?, + ?, + ?, + ?, + ?, + ?, + ? + )", + params![ + payload.scrobble.xata_id, + payload.scrobble.user_id.xata_id, + payload.scrobble.track_id.xata_id, + payload.scrobble.album_id.xata_id, + payload.scrobble.artist_id.xata_id, + payload.scrobble.uri, + payload.scrobble.xata_createdat, + ], + ) { + Ok(_) => (), + Err(e) => { + if !e.to_string().contains("violates primary key constraint") { + println!("[scrobbles] error: {}", e); + return Err(e.into()); + } + } + } + + Ok(()) +} + +pub async fn like(conn: Arc>, payload: LikePayload) -> Result<(), Error> { + let conn = conn.lock().unwrap(); + match conn.execute( + "INSERT INTO likes ( + id, + user_id, + track_id, + created_at + ) VALUES ( + ?, + ?, + ?, + ? + )", + params![ + payload.xata_id, + payload.user_id.xata_id, + payload.track_id.xata_id, + payload.xata_createdat, + ], + ) { + Ok(_) => (), + Err(e) => { + if !e.to_string().contains("violates primary key constraint") { + println!("[likes] error: {}", e); + return Err(e.into()); + } + } + } + Ok(()) +} + +pub async fn unlike(conn: Arc>, payload: UnlikePayload) -> Result<(), Error> { + let conn = conn.lock().unwrap(); + match conn.execute( + "DELETE FROM likes WHERE user_id = ? AND track_id = ?", + params![ + payload.user_id.xata_id, + payload.track_id.xata_id, + ], + ) { + Ok(_) => (), + Err(e) => { + println!("[unlikes] error: {}", e); + return Err(e.into()); + } + } + Ok(()) +} \ No newline at end of file diff --git a/crates/analytics/src/subscriber/types.rs b/crates/analytics/src/subscriber/types.rs new file mode 100644 index 00000000..e7611d97 --- /dev/null +++ b/crates/analytics/src/subscriber/types.rs @@ -0,0 +1,216 @@ +use chrono::{DateTime, Utc}; +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct LikePayload { + #[serde(skip_serializing_if = "Option::is_none")] + pub uri: Option, + pub track_id: Ref, + pub user_id: Ref, + pub xata_createdat: DateTime, + pub xata_id: String, + pub xata_updatedat: DateTime, + pub xata_version: i32, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct UnlikePayload { + #[serde(skip_serializing_if = "Option::is_none")] + pub uri: Option, + pub track_id: Ref, + pub user_id: Ref, + pub xata_createdat: DateTime, + pub xata_id: String, + pub xata_updatedat: DateTime, + pub xata_version: i32, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct ScrobblePayload { + pub scrobble: Scrobble, + pub user_album: UserAlbum, + pub user_artist: UserArtist, + pub user_track: UserTrack, + pub album_track: AlbumTrack, + pub artist_track: ArtistTrack, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct Scrobble { + pub album_id: AlbumId, + pub artist_id: ArtistId, + pub track_id: TrackId, + pub uri: String, + pub user_id: UserId, + pub xata_createdat: DateTime, + pub xata_id: String, + pub xata_updatedat: DateTime, + pub xata_version: i32, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct AlbumId { + #[serde(skip_serializing_if = "Option::is_none")] + pub album_art: Option, + pub artist: String, + pub artist_uri: String, + pub release_date: DateTime, + pub sha256: String, + pub title: String, + pub uri: String, + pub xata_createdat: DateTime, + pub xata_id: String, + pub xata_updatedat: DateTime, + pub xata_version: i32, + pub year: i32, + #[serde(skip_serializing_if = "Option::is_none")] + pub apple_music_link: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub spotify_link: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub tidal_link: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub youtube_link: Option, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct ArtistId { + pub name: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub picture: Option, + pub sha256: String, + pub uri: String, + pub xata_createdat: DateTime, + pub xata_id: String, + pub xata_updatedat: DateTime, + pub xata_version: i32, + #[serde(skip_serializing_if = "Option::is_none")] + pub apple_music_link: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub biography: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub born: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub born_in: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub died: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub spotify_link: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub tidal_link: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub youtube_link: Option, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct TrackId { + pub album: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub album_art: Option, + pub album_artist: String, + pub album_uri: String, + pub artist: String, + pub artist_uri: String, + pub disc_number: i32, + pub duration: i32, + pub sha256: String, + pub spotify_link: String, + pub title: String, + pub track_number: i32, + pub uri: String, + pub xata_createdat: DateTime, + pub xata_id: String, + pub xata_updatedat: DateTime, + pub xata_version: i32, + #[serde(skip_serializing_if = "Option::is_none")] + pub apple_music_link: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub composer: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub copyright_message: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub genre: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub label: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub lyrics: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub mb_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub tidal_link: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub youtube_link: Option, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct UserId { + pub avatar: String, + pub did: String, + pub display_name: String, + pub handle: String, + pub xata_createdat: DateTime, + pub xata_id: String, + pub xata_updatedat: DateTime, + pub xata_version: i32, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct UserAlbum { + pub album_id: Ref, + pub scrobbles: i32, + pub uri: String, + pub user_id: Ref, + pub xata_createdat: DateTime, + pub xata_id: String, + pub xata_updatedat: DateTime, + pub xata_version: i32, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct UserArtist { + pub artist_id: Ref, + pub scrobbles: i32, + pub uri: String, + pub user_id: Ref, + pub xata_createdat: DateTime, + pub xata_id: String, + pub xata_updatedat: DateTime, + pub xata_version: i32, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct UserTrack { + pub track_id: Ref, + pub scrobbles: i32, + pub uri: String, + pub user_id: Ref, + pub xata_createdat: DateTime, + pub xata_id: String, + pub xata_updatedat: DateTime, + pub xata_version: i32, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct AlbumTrack { + pub track_id: Ref, + pub album_id: Ref, + pub xata_createdat: DateTime, + pub xata_id: String, + pub xata_updatedat: DateTime, + pub xata_version: i32, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct ArtistTrack { + pub track_id: Ref, + pub artist_id: Ref, + pub xata_createdat: DateTime, + pub xata_id: String, + pub xata_updatedat: DateTime, + pub xata_version: i32, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct Ref { + pub xata_id: String, +} diff --git a/crates/analytics/src/types/album.rs b/crates/analytics/src/types/album.rs index ee81033b..d72e8100 100644 --- a/crates/analytics/src/types/album.rs +++ b/crates/analytics/src/types/album.rs @@ -43,3 +43,8 @@ pub struct GetTopAlbumsParams { pub user_did: Option, pub pagination: Option, } + +#[derive(Debug, Serialize, Deserialize)] +pub struct GetAlbumTracksParams { + pub album_id: String, +} diff --git a/crates/analytics/src/types/artist.rs b/crates/analytics/src/types/artist.rs index 12dee892..805983a7 100644 --- a/crates/analytics/src/types/artist.rs +++ b/crates/analytics/src/types/artist.rs @@ -44,3 +44,14 @@ pub struct GetTopArtistsParams { pub user_did: Option, pub pagination: Option, } + +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct GetArtistTracksParams { + pub artist_id: String, + pub pagination: Option, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct GetArtistAlbumsParams { + pub artist_id: String, +} diff --git a/crates/analytics/src/types/stats.rs b/crates/analytics/src/types/stats.rs index bd636f0e..fbd4edbb 100644 --- a/crates/analytics/src/types/stats.rs +++ b/crates/analytics/src/types/stats.rs @@ -68,3 +68,63 @@ impl Default for GetScrobblesPerYearParams { } } } + +#[derive(Debug, Serialize, Deserialize)] +pub struct GetAlbumScrobblesParams { + pub album_id: String, + pub start: Option, + pub end: Option, +} + +impl Default for GetAlbumScrobblesParams { + fn default() -> Self { + let current_date = Utc::now().naive_utc(); + let date_30_days_ago = current_date - Duration::days(30); + + GetAlbumScrobblesParams { + album_id: "".to_string(), + start: Some(date_30_days_ago.to_string()), + end: Some(current_date.to_string()), + } + } +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct GetArtistScrobblesParams { + pub artist_id: String, + pub start: Option, + pub end: Option, +} + +impl Default for GetArtistScrobblesParams { + fn default() -> Self { + let current_date = Utc::now().naive_utc(); + let date_30_days_ago = current_date - Duration::days(30); + + GetArtistScrobblesParams { + artist_id: "".to_string(), + start: Some(date_30_days_ago.to_string()), + end: Some(current_date.to_string()), + } + } +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct GetTrackScrobblesParams { + pub track_id: String, + pub start: Option, + pub end: Option, +} + +impl Default for GetTrackScrobblesParams { + fn default() -> Self { + let current_date = Utc::now().naive_utc(); + let date_30_days_ago = current_date - Duration::days(30); + + GetTrackScrobblesParams { + track_id: "".to_string(), + start: Some(date_30_days_ago.to_string()), + end: Some(current_date.to_string()), + } + } +} diff --git a/crates/analytics/src/xata/artist_album.rs b/crates/analytics/src/xata/artist_album.rs index 2e30ef11..89f6f1f7 100644 --- a/crates/analytics/src/xata/artist_album.rs +++ b/crates/analytics/src/xata/artist_album.rs @@ -5,4 +5,6 @@ pub struct ArtistAlbum { pub xata_id: String, pub artist_id: String, pub album_id: String, + #[serde(with = "chrono::serde::ts_seconds")] + pub xata_createdat: chrono::DateTime, } diff --git a/rockskyapi/rocksky-auth/.env.example b/rockskyapi/rocksky-auth/.env.example index 04d38391..8342b0e3 100644 --- a/rockskyapi/rocksky-auth/.env.example +++ b/rockskyapi/rocksky-auth/.env.example @@ -20,3 +20,6 @@ SPOTIFY_CLIENT_SECRET="your spotify" SPOTIFY_ENCRYPTION_KEY="your spotify encryption secret" SPOTIFY_ENCRYPRION_IV="your spotify encryption iv" ROCKSKY_BETA_TOKEN="your rocksky beta token" + +NATS_URL="nats://localhost:4222" +ANALYTICS_URL="http://localhost:7879" diff --git a/rockskyapi/rocksky-auth/src/context.ts b/rockskyapi/rocksky-auth/src/context.ts index 6bc761fd..84df9e9c 100644 --- a/rockskyapi/rocksky-auth/src/context.ts +++ b/rockskyapi/rocksky-auth/src/context.ts @@ -1,8 +1,10 @@ import { createClient } from "auth/client"; +import axios from "axios"; import { createDb, migrateToLatest } from "db"; import drizzle from "drizzle"; import { env } from "lib/env"; import { createBidirectionalResolver, createIdResolver } from "lib/idResolver"; +import { connect } from "nats"; import sqliteKv from "sqliteKv"; import { createStorage } from "unstorage"; import { getXataClient } from "xata"; @@ -25,6 +27,8 @@ export const ctx = { kv: new Map(), client, db: drizzle.db, + nc: await connect({ servers: env.NATS_URL }), + analytics: axios.create({ baseURL: env.ANALYTICS }), }; export type Context = typeof ctx; diff --git a/rockskyapi/rocksky-auth/src/lib/env.ts b/rockskyapi/rocksky-auth/src/lib/env.ts index 6e6fafe9..8dc2370b 100644 --- a/rockskyapi/rocksky-auth/src/lib/env.ts +++ b/rockskyapi/rocksky-auth/src/lib/env.ts @@ -25,4 +25,6 @@ export const env = cleanEnv(process.env, { SPOTIFY_ENCRYPTION_IV: str(), ROCKSKY_BETA_TOKEN: str({}), XATA_POSTGRES_URL: str({}), + NATS_URL: str({ devDefault: "nats://localhost:4222" }), + ANALYTICS: str({ devDefault: "http://localhost:7879" }), }); diff --git a/rockskyapi/rocksky-auth/src/lovedtracks/lovedtracks.service.ts b/rockskyapi/rocksky-auth/src/lovedtracks/lovedtracks.service.ts index e1a26110..2690ed80 100644 --- a/rockskyapi/rocksky-auth/src/lovedtracks/lovedtracks.service.ts +++ b/rockskyapi/rocksky-auth/src/lovedtracks/lovedtracks.service.ts @@ -137,6 +137,10 @@ export async function likeTrack(ctx: Context, track: Track, user) { track_id, } ); + + const message = JSON.stringify(created); + ctx.nc.publish("rocksky.like", Buffer.from(message)); + return created; } @@ -158,6 +162,9 @@ export async function unLikeTrack(ctx: Context, trackSha256: string, user) { return; } + const message = JSON.stringify(lovedTrack); + ctx.nc.publish("rocksky.unlike", Buffer.from(message)); + await ctx.client.db.loved_tracks.delete(lovedTrack.xata_id); } diff --git a/rockskyapi/rocksky-auth/src/nowplaying/nowplaying.service.ts b/rockskyapi/rocksky-auth/src/nowplaying/nowplaying.service.ts index 9163aa61..f9c4f079 100644 --- a/rockskyapi/rocksky-auth/src/nowplaying/nowplaying.service.ts +++ b/rockskyapi/rocksky-auth/src/nowplaying/nowplaying.service.ts @@ -317,6 +317,48 @@ export async function updateUserLibrary( } } +export async function publishScrobble(ctx: Context, id: string) { + const scrobble = await ctx.client.db.scrobbles + .select(["*", "track_id.*", "album_id.*", "artist_id.*", "user_id.*"]) + .filter("xata_id", equals(id)) + .getFirst(); + + const [user_album, user_artist, user_track, album_track, artist_track] = + await Promise.all([ + ctx.client.db.user_albums + .select(["*"]) + .filter("album_id.xata_id", equals(scrobble.album_id.xata_id)) + .getFirst(), + ctx.client.db.user_artists + .select(["*"]) + .filter("artist_id.xata_id", equals(scrobble.artist_id.xata_id)) + .getFirst(), + ctx.client.db.user_tracks + .select(["*"]) + .filter("track_id.xata_id", equals(scrobble.track_id.xata_id)) + .getFirst(), + ctx.client.db.album_tracks + .select(["*"]) + .filter("track_id.xata_id", equals(scrobble.track_id.xata_id)) + .getFirst(), + ctx.client.db.artist_tracks + .select(["*"]) + .filter("track_id.xata_id", equals(scrobble.track_id.xata_id)) + .getFirst(), + ]); + + const message = JSON.stringify({ + scrobble, + user_album, + user_artist, + user_track, + album_track, + artist_track, + }); + + ctx.nc.publish("rocksky.scrobble", Buffer.from(message)); +} + export async function scrobbleTrack( ctx: Context, track: Track, @@ -491,5 +533,7 @@ export async function scrobbleTrack( uri: scrobbleUri, }); + await publishScrobble(ctx, scrobble.xata_id); + return scrobble; } diff --git a/rockskyapi/rocksky-auth/src/users/app.ts b/rockskyapi/rocksky-auth/src/users/app.ts index a8ec0419..fca372b9 100644 --- a/rockskyapi/rocksky-auth/src/users/app.ts +++ b/rockskyapi/rocksky-auth/src/users/app.ts @@ -51,27 +51,15 @@ app.get("/:handle/scrobbles", async (c) => { const size = +c.req.query("size") || 10; const offset = +c.req.query("offset") || 0; - const scrobbles = await ctx.client.db.scrobbles - .select(["track_id.*", "uri", "album_id.*", "artist_id.*"]) - .filter({ - $any: [ - { - "user_id.did": handle, - }, - { - "user_id.handle": handle, - }, - ], - }) - .sort("xata_createdat", "desc") - .getPaginated({ - pagination: { - size, - offset, - }, - }); + const { data } = await ctx.analytics.post("library.getScrobbles", { + user_did: handle, + pagination: { + skip: offset, + take: size, + }, + }); - return c.json(scrobbles.records); + return c.json(data); }); app.get("/:did/albums", async (c) => { @@ -79,29 +67,15 @@ app.get("/:did/albums", async (c) => { const size = +c.req.query("size") || 10; const offset = +c.req.query("offset") || 0; - const albums = await ctx.client.db.user_albums - .select(["album_id.*", "scrobbles"]) - .filter({ - $any: [ - { - "user_id.did": did, - }, - { - "user_id.handle": did, - }, - ], - }) - .getPaginated({ - sort: { - scrobbles: "desc", - }, - pagination: { - size, - offset, - }, - }); + const { data } = await ctx.analytics.post("library.getTopAlbums", { + user_did: did, + pagination: { + skip: offset, + take: size, + }, + }); - return c.json(albums.records.map((item) => ({ ...item.album_id, tags: [] }))); + return c.json(data.map((item) => ({ ...item, tags: [] }))); }); app.get("/:did/artists", async (c) => { @@ -109,31 +83,15 @@ app.get("/:did/artists", async (c) => { const size = +c.req.query("size") || 10; const offset = +c.req.query("offset") || 0; - const artists = await ctx.client.db.user_artists - .select(["artist_id.*", "scrobbles"]) - .filter({ - $any: [ - { - "user_id.did": did, - }, - { - "user_id.handle": did, - }, - ], - }) - .getPaginated({ - sort: { - scrobbles: "desc", - }, - pagination: { - size, - offset, - }, - }); + const { data } = await ctx.analytics.post("library.getTopArtists", { + user_did: did, + pagination: { + skip: offset, + take: size, + }, + }); - return c.json( - artists.records.map((item) => ({ ...item.artist_id, tags: [] })) - ); + return c.json(data.map((item) => ({ ...item, tags: [] }))); }); app.get("/:did/tracks", async (c) => { @@ -141,32 +99,17 @@ app.get("/:did/tracks", async (c) => { const size = +c.req.query("size") || 10; const offset = +c.req.query("offset") || 0; - const artists = await ctx.client.db.user_tracks - .select(["track_id.*", "scrobbles"]) - .filter({ - $any: [ - { - "user_id.did": did, - }, - { - "user_id.handle": did, - }, - ], - }) - .getPaginated({ - sort: { - scrobbles: "desc", - }, - pagination: { - size, - offset, - }, - }); + const { data } = await ctx.analytics.post("library.getTopTracks", { + user_did: did, + pagination: { + skip: offset, + take: size, + }, + }); return c.json( - artists.records.map((item) => ({ - ...item.track_id, - scrobles: item.scrobbles, + data.map((item) => ({ + ...item, tags: [], })) ); @@ -1202,58 +1145,15 @@ app.get("/:did/app.rocksky.shout/:rkey/replies", async (c) => { app.get("/:did/stats", async (c) => { const did = c.req.param("did"); - const scrobbles = await ctx.client.db.scrobbles - .select(["user_id.*"]) - .filter({ - $any: [ - { - "user_id.did": did, - }, - { - "user_id.handle": did, - }, - ], - }) - .summarize({ - summaries: { - total: { - count: "*", - }, - }, - }); - - const artists = await ctx.client.db.user_artists - .select(["artist_id.*", "user_id.*"]) - .filter({ - $any: [ - { - "user_id.did": did, - }, - { - "user_id.handle": did, - }, - ], - }) - .getAll(); - const lovedTracks = await ctx.client.db.loved_tracks - .select(["track_id.*", "user_id.*"]) - .filter({ - $any: [ - { - "user_id.did": did, - }, - { - "user_id.handle": did, - }, - ], - }) - .getAll(); + const { data } = await ctx.analytics.post("library.getStats", { + user_did: did, + }); return c.json({ - scrobbles: _.get(scrobbles, "summaries.0.total", 0), - artists: artists.length, - lovedTracks: lovedTracks.length, + scrobbles: data.scrobbles, + artists: data.artists, + lovedTracks: data.loved_tracks, }); }); diff --git a/rockskyweb/src/pages/profile/overview/recenttracks/RecentTracks.tsx b/rockskyweb/src/pages/profile/overview/recenttracks/RecentTracks.tsx index 9bae5d4f..3633a3ca 100644 --- a/rockskyweb/src/pages/profile/overview/recenttracks/RecentTracks.tsx +++ b/rockskyweb/src/pages/profile/overview/recenttracks/RecentTracks.tsx @@ -56,19 +56,18 @@ function RecentTracks(props: RecentTracksProps) { const getRecentTracks = async () => { const data = await getRecentTracksByDid(did, 0, props.size); setRecentTracks( - data.map(({ track_id, album_id, artist_id, uri, xata_createdat }) => ({ - id: track_id.xata_id, - title: track_id.title, - artist: track_id.artist, - album: track_id.album, - albumArt: track_id.album_art, - albumArtist: track_id.album_artist, - duration: track_id.duration, - uri: track_id.uri, - date: xata_createdat, - scrobbleUri: uri, - albumUri: album_id.uri, - artistUri: artist_id.uri, + data.map((item) => ({ + id: item.id, + title: item.title, + artist: item.artist, + album: item.album, + albumArt: item.album_art, + albumArtist: item.album_artist, + uri: item.uri, + date: item.created_at, + scrobbleUri: item.uri, + albumUri: item.album_uri, + artistUri: item.artist_uri, })) ); }; diff --git a/rockskyweb/src/types/scrobble.ts b/rockskyweb/src/types/scrobble.ts index 518535e6..72db8b33 100644 --- a/rockskyweb/src/types/scrobble.ts +++ b/rockskyweb/src/types/scrobble.ts @@ -1,31 +1,15 @@ export type Scrobble = { - track_id: { - xata_id: string; - title: string; - artist: string; - album: string; - album_art?: string; - album_artist: string; - uri: string; - duration: number; - }; - album_id: { - uri: string; - album_art?: string; - artist: string; - release_date: string; - title: string; - xata_id: string; - year: number; - }; - artist_id: { - uri: string; - name: string; - xata_id: string; - bio?: string; - picture?: string; - }; + id: string; + track_id: string; + title: string; + artist: string; + album: string; + album_art?: string; + album_artist: string; + handle: string; + track_uri: string; + album_uri: string; + artist_uri: string; uri: string; - xata_createdat: string; - xata_id: string; + created_at: string; };