From 73aaf312cbfd165b9951ea1c81e1979aab194a4c Mon Sep 17 00:00:00 2001 From: Tsiry Sandratraina Date: Sun, 2 Mar 2025 11:21:58 +0300 Subject: [PATCH] finish analytics API implementation --- Cargo.lock | 1 + crates/analytics/Cargo.toml | 1 + crates/analytics/src/cmd/serve.rs | 44 ++- crates/analytics/src/cmd/sync.rs | 26 +- crates/analytics/src/core.rs | 35 +- crates/analytics/src/handlers/albums.rs | 187 ++++++++++ crates/analytics/src/handlers/artists.rs | 199 +++++++++++ crates/analytics/src/handlers/mod.rs | 50 +++ crates/analytics/src/handlers/scrobbles.rs | 117 +++++++ crates/analytics/src/handlers/stats.rs | 222 ++++++++++++ crates/analytics/src/handlers/tracks.rs | 344 +++++++++++++++++++ crates/analytics/src/main.rs | 12 +- crates/analytics/src/types/album.rs | 37 +- crates/analytics/src/types/album_track.rs | 0 crates/analytics/src/types/artist.rs | 37 +- crates/analytics/src/types/artist_album.rs | 0 crates/analytics/src/types/artist_track.rs | 0 crates/analytics/src/types/filters.rs | 8 + crates/analytics/src/types/loved_track.rs | 0 crates/analytics/src/types/mod.rs | 9 +- crates/analytics/src/types/pagination.rs | 16 + crates/analytics/src/types/playlist.rs | 17 +- crates/analytics/src/types/playlist_track.rs | 0 crates/analytics/src/types/scrobble.rs | 47 ++- crates/analytics/src/types/stats.rs | 70 ++++ crates/analytics/src/types/track.rs | 50 ++- crates/analytics/src/types/user.rs | 16 +- crates/analytics/src/types/user_album.rs | 0 crates/analytics/src/types/user_artist.rs | 0 crates/analytics/src/types/user_playlist.rs | 0 crates/analytics/src/types/user_track.rs | 0 31 files changed, 1434 insertions(+), 111 deletions(-) create mode 100644 crates/analytics/src/handlers/albums.rs create mode 100644 crates/analytics/src/handlers/artists.rs create mode 100644 crates/analytics/src/handlers/mod.rs create mode 100644 crates/analytics/src/handlers/scrobbles.rs create mode 100644 crates/analytics/src/handlers/stats.rs create mode 100644 crates/analytics/src/handlers/tracks.rs delete mode 100644 crates/analytics/src/types/album_track.rs delete mode 100644 crates/analytics/src/types/artist_album.rs delete mode 100644 crates/analytics/src/types/artist_track.rs create mode 100644 crates/analytics/src/types/filters.rs delete mode 100644 crates/analytics/src/types/loved_track.rs create mode 100644 crates/analytics/src/types/pagination.rs delete mode 100644 crates/analytics/src/types/playlist_track.rs create mode 100644 crates/analytics/src/types/stats.rs delete mode 100644 crates/analytics/src/types/user_album.rs delete mode 100644 crates/analytics/src/types/user_artist.rs delete mode 100644 crates/analytics/src/types/user_playlist.rs delete mode 100644 crates/analytics/src/types/user_track.rs diff --git a/Cargo.lock b/Cargo.lock index aa7a02c2..4ae0127a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -277,6 +277,7 @@ dependencies = [ "clap", "dotenv", "duckdb", + "futures-util", "owo-colors", "polars", "serde", diff --git a/crates/analytics/Cargo.toml b/crates/analytics/Cargo.toml index 5cac9a6e..eb3650c7 100644 --- a/crates/analytics/Cargo.toml +++ b/crates/analytics/Cargo.toml @@ -27,3 +27,4 @@ anyhow = "1.0.96" polars = "0.46.0" clap = "4.5.31" actix-web = "4.9.0" +futures-util = "0.3.31" diff --git a/crates/analytics/src/cmd/serve.rs b/crates/analytics/src/cmd/serve.rs index 99d99ec6..7e35cd0b 100644 --- a/crates/analytics/src/cmd/serve.rs +++ b/crates/analytics/src/cmd/serve.rs @@ -1,27 +1,51 @@ use std::env; -use actix_web::{get, App, HttpRequest, HttpServer}; +use actix_web::{get, post, web::{self, Data}, App, HttpRequest, HttpServer, Responder}; use duckdb::Connection; use anyhow::Error; use owo_colors::OwoColorize; +use std::sync::{Arc, Mutex}; + +use crate::handlers::handle; #[get("/")] async fn index(_req: HttpRequest) -> String { "Hello world!".to_owned() } -pub async fn serve(_conn: &Connection) -> Result<(), Error> { - let host = env::var("HOST").unwrap_or_else(|_| "127.0.0.1".to_string()); - let port = env::var("PORT").unwrap_or_else(|_| "7879".to_string()); +#[post("/{method}")] +async fn call_method( + data: web::Data>>, + mut payload: web::Payload, + req: HttpRequest) -> Result { + let method = req.match_info().get("method").unwrap_or("unknown"); + println!("Method: {}", method.bright_green()); + + let conn = data.get_ref().clone(); + handle(method, &mut payload, &req, conn).await + .map_err(actix_web::error::ErrorInternalServerError) +} + + +pub async fn serve(conn: Arc>) -> Result<(), Error> { + 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); let url = format!("http://{}", addr); println!("Listening on {}", url.bright_green()); - HttpServer::new(|| App::new().service(index)) - .bind(&addr)? - .run() - .await - .map_err(Error::new)?; + let conn = conn.clone(); + HttpServer::new(move || { + App::new() + .app_data(Data::new( conn.clone())) + .service(index) + .service(call_method) + }) + .bind(&addr)? + .run() + .await + .map_err(Error::new)?; + Ok(()) -} \ No newline at end of file +} diff --git a/crates/analytics/src/cmd/sync.rs b/crates/analytics/src/cmd/sync.rs index 9ff53e63..88569032 100644 --- a/crates/analytics/src/cmd/sync.rs +++ b/crates/analytics/src/cmd/sync.rs @@ -1,20 +1,22 @@ +use std::sync::{Arc, Mutex}; + use anyhow::Error; use duckdb::Connection; use sqlx::{Pool, Postgres}; use crate::core::*; -pub async fn sync(conn: &Connection, pool: &Pool) -> Result<(), Error> { - load_tracks(conn, pool).await?; - load_artists(conn, pool).await?; - load_albums(conn, pool).await?; - load_users(conn, pool).await?; - load_scrobbles(conn, pool).await?; - load_album_tracks(conn, pool).await?; - load_loved_tracks(conn, pool).await?; - load_artist_tracks(conn, pool).await?; - load_user_albums(conn, pool).await?; - load_user_artists(conn, pool).await?; - load_user_tracks(conn, pool).await?; +pub async fn sync(conn: Arc>, pool: &Pool) -> Result<(), Error> { + load_tracks(conn.clone(), pool).await?; + load_artists(conn.clone(), pool).await?; + load_albums(conn.clone(), pool).await?; + load_users(conn.clone(), pool).await?; + load_scrobbles(conn.clone(), pool).await?; + load_album_tracks(conn.clone(), pool).await?; + load_loved_tracks(conn.clone(), pool).await?; + load_artist_tracks(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?; Ok(()) } diff --git a/crates/analytics/src/core.rs b/crates/analytics/src/core.rs index 4386920b..1f1c475e 100644 --- a/crates/analytics/src/core.rs +++ b/crates/analytics/src/core.rs @@ -1,3 +1,5 @@ +use std::sync::{Arc, Mutex}; + use duckdb::{params, Connection}; use anyhow::Error; use owo_colors::OwoColorize; @@ -174,7 +176,8 @@ pub async fn create_tables(conn: &Connection) -> Result<(), Error> { Ok(()) } -pub async fn load_tracks(conn: &Connection, pool: &Pool) -> Result<(), Error> { +pub async fn load_tracks(conn: Arc>, pool: &Pool) -> Result<(), Error> { + let conn = conn.lock().unwrap(); let tracks: Vec = sqlx::query_as(r#" SELECT * FROM tracks "#) @@ -246,7 +249,8 @@ pub async fn load_tracks(conn: &Connection, pool: &Pool) -> Result<(), Ok(()) } -pub async fn load_artists(conn: &Connection, pool: &Pool) -> Result<(), Error> { +pub async fn load_artists(conn: Arc>, pool: &Pool) -> Result<(), Error> { + let conn = conn.lock().unwrap(); let artists: Vec = sqlx::query_as(r#" SELECT * FROM artists "#) @@ -308,7 +312,8 @@ pub async fn load_artists(conn: &Connection, pool: &Pool) -> Result<() Ok(()) } -pub async fn load_albums(conn: &Connection, pool: &Pool) -> Result<(), Error> { +pub async fn load_albums(conn: Arc>, pool: &Pool) -> Result<(), Error> { + let conn = conn.lock().unwrap(); let albums: Vec = sqlx::query_as(r#" SELECT * FROM albums "#) @@ -370,7 +375,8 @@ pub async fn load_albums(conn: &Connection, pool: &Pool) -> Result<(), Ok(()) } -pub async fn load_users(conn: &Connection, pool: &Pool) -> Result<(), Error> { +pub async fn load_users(conn: Arc>, pool: &Pool) -> Result<(), Error> { + let conn = conn.lock().unwrap(); let users: Vec = sqlx::query_as(r#" SELECT * FROM users "#) @@ -408,7 +414,8 @@ pub async fn load_users(conn: &Connection, pool: &Pool) -> Result<(), Ok(()) } -pub async fn load_scrobbles(conn: &Connection, pool: &Pool) -> Result<(), Error> { +pub async fn load_scrobbles(conn: Arc>, pool: &Pool) -> Result<(), Error> { + let conn = conn.lock().unwrap(); let scrobbles: Vec = sqlx::query_as(r#" SELECT * FROM scrobbles "#) @@ -457,7 +464,8 @@ pub async fn load_scrobbles(conn: &Connection, pool: &Pool) -> Result< Ok(()) } -pub async fn load_album_tracks(conn: &Connection, pool: &Pool) -> Result<(), Error> { +pub async fn load_album_tracks(conn: Arc>, pool: &Pool) -> Result<(), Error> { + let conn = conn.lock().unwrap(); let album_tracks: Vec = sqlx::query_as(r#" SELECT * FROM album_tracks "#) @@ -488,7 +496,8 @@ pub async fn load_album_tracks(conn: &Connection, pool: &Pool) -> Resu Ok(()) } -pub async fn load_loved_tracks(conn: &Connection, pool: &Pool) -> Result<(), Error> { +pub async fn load_loved_tracks(conn: Arc>, pool: &Pool) -> Result<(), Error> { + let conn = conn.lock().unwrap(); let loved_tracks: Vec = sqlx::query_as(r#" SELECT * FROM loved_tracks "#) @@ -523,7 +532,8 @@ pub async fn load_loved_tracks(conn: &Connection, pool: &Pool) -> Resu Ok(()) } -pub async fn load_artist_tracks(conn: &Connection, pool: &Pool) -> Result<(), Error> { +pub async fn load_artist_tracks(conn: Arc>, pool: &Pool) -> Result<(), Error> { + let conn = conn.lock().unwrap(); let artist_tracks: Vec = sqlx::query_as(r#" SELECT * FROM artist_tracks "#) @@ -550,7 +560,8 @@ pub async fn load_artist_tracks(conn: &Connection, pool: &Pool) -> Res Ok(()) } -pub async fn load_user_albums(conn: &Connection, pool: &Pool) -> Result<(), Error> { +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#" SELECT * FROM user_albums "#) @@ -577,7 +588,8 @@ pub async fn load_user_albums(conn: &Connection, pool: &Pool) -> Resul Ok(()) } -pub async fn load_user_artists(conn: &Connection, pool: &Pool) -> Result<(), Error> { +pub async fn load_user_artists(conn: Arc>, pool: &Pool) -> Result<(), Error> { + let conn = conn.lock().unwrap(); let user_artists: Vec = sqlx::query_as(r#" SELECT * FROM user_artists "#) @@ -604,7 +616,8 @@ pub async fn load_user_artists(conn: &Connection, pool: &Pool) -> Resu Ok(()) } -pub async fn load_user_tracks(conn: &Connection, pool: &Pool) -> Result<(), Error> { +pub async fn load_user_tracks(conn: Arc>, pool: &Pool) -> Result<(), Error> { + let conn = conn.lock().unwrap(); let user_tracks: Vec = sqlx::query_as(r#" SELECT * FROM user_tracks "#) diff --git a/crates/analytics/src/handlers/albums.rs b/crates/analytics/src/handlers/albums.rs new file mode 100644 index 00000000..571ab49e --- /dev/null +++ b/crates/analytics/src/handlers/albums.rs @@ -0,0 +1,187 @@ +use std::sync::{Arc, Mutex}; + +use actix_web::{web, HttpRequest, HttpResponse}; +use analytics::types::album::{Album, GetAlbumsParams, GetTopAlbumsParams}; +use duckdb::Connection; +use anyhow::Error; +use futures_util::StreamExt; + +use crate::read_payload; + +pub async fn get_albums(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 did = params.user_did; + + let conn = conn.lock().unwrap(); + let mut stmt = match did { + Some(_) => { + conn.prepare(r#" + 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 = ? + ORDER BY a.title ASC OFFSET ? LIMIT ?; + "#)? + }, + None => { + conn.prepare("SELECT * FROM albums ORDER BY title ASC OFFSET ? LIMIT ?")? + } + }; + + match did { + Some(did) => { + let albums_iter = stmt.query_map([did, limit.to_string(), offset.to_string()], |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: row.get(6)?, + tidal_link: row.get(7)?, + youtube_link: row.get(8)?, + apple_music_link: row.get(9)?, + sha256: row.get(10)?, + uri: row.get(11)?, + artist_uri: row.get(12)?, + ..Default::default() + }) + })?; + + let albums: Result, _> = albums_iter.collect(); + Ok(HttpResponse::Ok().json(web::Json(albums?))) + }, + None => { + let albums_iter = stmt.query_map([limit, offset], |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: row.get(6)?, + tidal_link: row.get(7)?, + youtube_link: row.get(8)?, + apple_music_link: row.get(9)?, + sha256: row.get(10)?, + uri: row.get(11)?, + artist_uri: row.get(12)?, + ..Default::default() + }) + })?; + + let albums: Result, _> = albums_iter.collect(); + Ok(HttpResponse::Ok().json(web::Json(albums?))) + } + } +} + + +pub async fn get_top_albums(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 did = params.user_did; + + let conn = conn.lock().unwrap(); + let mut stmt = match did { + Some(_) => conn.prepare(r#" + SELECT + s.album_id AS id, + a.title AS title, + ar.name AS artist, + a.album_art AS album_art, + a.release_date, + a.year, + a.uri AS uri, + COUNT(*) AS play_count, + COUNT(DISTINCT s.user_id) AS unique_listeners + FROM + scrobbles s + LEFT JOIN + albums a ON s.album_id = a.id + LEFT JOIN + 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 = ? + GROUP BY + s.album_id, a.title, ar.name, a.release_date, a.year, a.uri, a.album_art + ORDER BY + play_count DESC + OFFSET ? + LIMIT ?; + "#)?, + None => conn.prepare(r#" + SELECT + s.album_id AS id, + a.title AS title, + ar.name AS artist, + a.album_art AS album_art, + a.release_date, + a.year, + a.uri AS uri, + COUNT(*) AS play_count, + COUNT(DISTINCT s.user_id) AS unique_listeners + FROM + scrobbles s + LEFT JOIN + albums a ON s.album_id = a.id + 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 + ORDER BY + play_count DESC + OFFSET ? + LIMIT ?; + "#)? + }; + + match did { + Some(did) => { + let albums = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + Ok(Album { + id: row.get(0)?, + title: row.get(1)?, + artist: row.get(2)?, + album_art: row.get(3)?, + 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)?), + ..Default::default() + }) + })?; + let albums: Result, _> = albums.collect(); + Ok(HttpResponse::Ok().json(web::Json(albums?))) + }, + None => { + let albums = stmt.query_map([limit, offset], |row| { + Ok(Album { + id: row.get(0)?, + title: row.get(1)?, + artist: row.get(2)?, + album_art: row.get(3)?, + 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)?), + ..Default::default() + }) + })?; + let albums: Result, _> = albums.collect(); + Ok(HttpResponse::Ok().json(web::Json(albums?))) + } + } +} \ No newline at end of file diff --git a/crates/analytics/src/handlers/artists.rs b/crates/analytics/src/handlers/artists.rs new file mode 100644 index 00000000..db0cfa6c --- /dev/null +++ b/crates/analytics/src/handlers/artists.rs @@ -0,0 +1,199 @@ +use std::sync::{Arc, Mutex}; + +use actix_web::{web, HttpRequest, HttpResponse}; +use analytics::types::artist::{Artist, GetTopArtistsParams}; +use duckdb::Connection; +use anyhow::Error; +use futures_util::StreamExt; + +use crate::{read_payload, types::artist::GetArtistsParams}; + +pub async fn get_artists(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 did = params.user_did; + + let conn = conn.lock().unwrap(); + let mut stmt = match did { + Some(_) => { + conn.prepare(r#" + 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 = ? + ORDER BY a.name ASC OFFSET ? LIMIT ?; + "#)? + }, + None => { + conn.prepare("SELECT * FROM artists ORDER BY name ASC OFFSET ? LIMIT ?")? + } + }; + + match did { + Some(did) => { + let artists = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + Ok(Artist { + id: row.get(0)?, + name: row.get(1)?, + biography: row.get(2)?, + born: row.get(3)?, + born_in: row.get(4)?, + died: row.get(5)?, + picture: row.get(6)?, + sha256: row.get(7)?, + spotify_link: row.get(8)?, + tidal_link: row.get(9)?, + youtube_link: row.get(10)?, + apple_music_link: row.get(11)?, + uri: row.get(12)?, + play_count: None, + unique_listeners: None, + }) + })?; + + let artists: Result, _> = artists.collect(); + Ok(HttpResponse::Ok().json(artists?)) + }, + None => { + let artists = stmt.query_map([limit, offset], |row| { + Ok(Artist { + id: row.get(0)?, + name: row.get(1)?, + biography: row.get(2)?, + born: row.get(3)?, + born_in: row.get(4)?, + died: row.get(5)?, + picture: row.get(6)?, + sha256: row.get(7)?, + spotify_link: row.get(8)?, + tidal_link: row.get(9)?, + youtube_link: row.get(10)?, + apple_music_link: row.get(11)?, + uri: row.get(12)?, + play_count: None, + unique_listeners: None, + }) + })?; + + let artists: Result, _> = artists.collect(); + Ok(HttpResponse::Ok().json(artists?)) + } + } +} + +pub async fn get_top_artists(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 did = params.user_did; + + let conn = conn.lock().unwrap(); + let mut stmt = match did { + Some(_) => { + conn.prepare(r#" + SELECT + s.artist_id AS id, + ar.name AS artist_name, + ar.picture AS picture, + ar.sha256 AS sha256, + ar.uri AS uri, + COUNT(*) AS play_count, + COUNT(DISTINCT s.user_id) AS unique_listeners + FROM + scrobbles s + LEFT JOIN + artists ar ON s.artist_id = ar.id + LEFT JOIN + users u ON s.user_id = u.id + WHERE + s.artist_id IS NOT NULL AND u.did = ? + GROUP BY + s.artist_id, ar.name, ar.uri, ar.picture, ar.sha256 + ORDER BY + play_count DESC + OFFSET ? + LIMIT ?; + "#)? + }, + None => { + conn.prepare(r#" + SELECT + s.artist_id AS id, + ar.name AS artist_name, + ar.picture AS picture, + ar.sha256 AS sha256, + ar.uri AS uri, + COUNT(*) AS play_count, + COUNT(DISTINCT s.user_id) AS unique_listeners + FROM + scrobbles s + LEFT JOIN + artists ar ON s.artist_id = ar.id + WHERE + s.artist_id IS NOT NULL + GROUP BY + s.artist_id, ar.name, ar.uri, ar.picture, ar.sha256 + ORDER BY + play_count DESC + OFFSET ? + LIMIT ?; + "#)? + } + }; + + match did { + Some(did) => { + let artists = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + Ok(Artist { + id: row.get(0)?, + name: row.get(1)?, + biography: None, + born: None, + born_in: None, + died: None, + picture: row.get(2)?, + sha256: row.get(3)?, + spotify_link: None, + tidal_link: None, + youtube_link: None, + apple_music_link: None, + uri: row.get(4)?, + play_count: Some(row.get(5)?), + unique_listeners: Some(row.get(6)?), + }) + })?; + + let artists: Result, _> = artists.collect(); + Ok(HttpResponse::Ok().json(artists?)) + }, + None => { + let artists = stmt.query_map([limit, offset], |row| { + Ok(Artist { + id: row.get(0)?, + name: row.get(1)?, + biography: None, + born: None, + born_in: None, + died: None, + picture: row.get(2)?, + sha256: row.get(3)?, + spotify_link: None, + tidal_link: None, + youtube_link: None, + apple_music_link: None, + uri: row.get(4)?, + play_count: Some(row.get(5)?), + unique_listeners: Some(row.get(6)?), + }) + })?; + + let artists: Result, _> = artists.collect(); + Ok(HttpResponse::Ok().json(artists?)) + } + } +} \ No newline at end of file diff --git a/crates/analytics/src/handlers/mod.rs b/crates/analytics/src/handlers/mod.rs new file mode 100644 index 00000000..f104148e --- /dev/null +++ b/crates/analytics/src/handlers/mod.rs @@ -0,0 +1,50 @@ +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 duckdb::Connection; +use scrobbles::get_scrobbles; +use stats::{get_scrobbles_per_day, get_scrobbles_per_month, get_scrobbles_per_year, get_stats}; +use tracks::{get_loved_tracks, get_top_tracks, get_tracks}; +use anyhow::Error; + +pub mod albums; +pub mod artists; +pub mod scrobbles; +pub mod tracks; +pub mod stats; + + +#[macro_export] +macro_rules! read_payload { + ($payload:expr) => {{ + let mut body = Vec::new(); + while let Some(chunk) = $payload.next().await { + match chunk { + Ok(bytes) => body.extend_from_slice(&bytes), + Err(err) => return Err(err.into()), + } + } + body + }}; +} + + +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, + "library.getArtists" => get_artists(payload, req, conn.clone()).await, + "library.getTracks" => get_tracks(payload, req, conn.clone()).await, + "library.getScrobbles" => get_scrobbles(payload, req, conn.clone()).await, + "library.getLovedTracks" => get_loved_tracks(payload, req, conn.clone()).await, + "library.getStats" => get_stats(payload, req, conn.clone()).await, + "library.getTopAlbums" => get_top_albums(payload, req, conn.clone()).await, + "library.getTopArtists" => get_top_artists(payload, req, conn.clone()).await, + "library.getTopTracks" => get_top_tracks(payload, req, conn.clone()).await, + "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, + _ => return Err(anyhow::anyhow!("Method not found")), + } +} diff --git a/crates/analytics/src/handlers/scrobbles.rs b/crates/analytics/src/handlers/scrobbles.rs new file mode 100644 index 00000000..f12542df --- /dev/null +++ b/crates/analytics/src/handlers/scrobbles.rs @@ -0,0 +1,117 @@ +use std::sync::{Arc, Mutex}; + +use actix_web::{web, HttpRequest, HttpResponse}; +use analytics::types::scrobble::{GetScrobblesParams, ScrobbleTrack}; +use duckdb::Connection; +use anyhow::Error; +use futures_util::StreamExt; + +use crate::read_payload; + +pub async fn get_scrobbles(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 did = params.user_did; + + let conn = conn.lock().unwrap(); + let mut stmt = match did { + Some(_) => conn.prepare(r#" + SELECT + s.id, + t.id as track_id, + t.title, + t.artist, + t.album_artist, + t.album, + t.album_art, + u.handle, + s.uri, + t.uri as track_uri, + a.uri as artist_uri, + al.uri as album_uri, + s.created_at + FROM scrobbles s + LEFT JOIN artists a ON s.artist_id = a.id + 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 = ? + 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 ? + LIMIT ?; + "#)?, + None => conn.prepare(r#" + SELECT + s.id, + t.id as track_id, + t.title, + t.artist, + t.album_artist, + t.album, + t.album_art, + u.handle, + s.uri, + t.uri as track_uri, + a.uri as artist_uri, + al.uri as album_uri, + s.created_at + FROM scrobbles s + LEFT JOIN artists a ON s.artist_id = a.id + 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 + 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 ? + LIMIT ?; + "#)?, + }; + match did { + Some(did) => { + let scrobbles = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + Ok(ScrobbleTrack { + id: row.get(0)?, + track_id: row.get(1)?, + title: row.get(2)?, + artist: row.get(3)?, + album_artist: row.get(4)?, + album: row.get(5)?, + album_art: row.get(6)?, + handle: row.get(7)?, + uri: row.get(8)?, + track_uri: row.get(9)?, + artist_uri: row.get(10)?, + album_uri: row.get(11)?, + created_at: row.get(12)?, + }) + })?; + let scrobbles: Result, _> = scrobbles.collect(); + Ok(HttpResponse::Ok().json(scrobbles?)) + }, + None => { + let scrobbles = stmt.query_map([limit, offset], |row| { + Ok(ScrobbleTrack { + id: row.get(0)?, + track_id: row.get(1)?, + title: row.get(2)?, + artist: row.get(3)?, + album_artist: row.get(4)?, + album: row.get(5)?, + album_art: row.get(6)?, + handle: row.get(7)?, + uri: row.get(8)?, + track_uri: row.get(9)?, + artist_uri: row.get(10)?, + album_uri: row.get(11)?, + created_at: row.get(12)?, + }) + })?; + let scrobbles: Result, _> = scrobbles.collect(); + Ok(HttpResponse::Ok().json(scrobbles?)) + } + } +} diff --git a/crates/analytics/src/handlers/stats.rs b/crates/analytics/src/handlers/stats.rs new file mode 100644 index 00000000..76bc6a51 --- /dev/null +++ b/crates/analytics/src/handlers/stats.rs @@ -0,0 +1,222 @@ +use std::sync::{Arc, Mutex}; + +use actix_web::{web, HttpRequest, HttpResponse}; +use analytics::types::{scrobble::{ScrobblesPerDay, ScrobblesPerMonth, ScrobblesPerYear}, stats::{GetScrobblesPerDayParams, GetScrobblesPerMonthParams, GetScrobblesPerYearParams, GetStatsParams}}; +use duckdb::Connection; +use anyhow::Error; +use serde_json::json; +use futures_util::StreamExt; +use crate::read_payload; + +pub async fn get_stats(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("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 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 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 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_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))?; + + Ok(HttpResponse::Ok().json(json!({ + "scrobbles": scrobbles, + "artists": artists, + "loved_tracks": loved_tracks, + "albums": albums, + "tracks": tracks, + }))) +} + +pub async fn get_scrobbles_per_day(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(GetScrobblesPerDayParams::default().start.unwrap()); + let end = params.end.unwrap_or(GetScrobblesPerDayParams::default().end.unwrap()); + let did = params.user_did; + + let conn = conn.lock().unwrap(); + match did { + Some(did) => { + let mut stmt = conn.prepare(r#" + SELECT + date_trunc('day', created_at) AS date, + COUNT(track_id) AS count + FROM + scrobbles + LEFT JOIN users u ON scrobbles.user_id = u.id + WHERE + u.did = ? + AND created_at BETWEEN ? AND ? + GROUP BY + date_trunc('day', created_at) + ORDER BY + date; + "#)?; + let scrobbles = stmt.query_map([did, start, end], |row| { + Ok(ScrobblesPerDay { + date: row.get(0)?, + count: row.get(1)?, + }) + })?; + let scrobbles: Result, _> = scrobbles.collect(); + Ok(HttpResponse::Ok().json(scrobbles?)) + }, + None => { + let mut stmt = conn.prepare(r#" + SELECT + date_trunc('day', created_at) AS date, + COUNT(track_id) AS count + FROM + scrobbles + WHERE + created_at BETWEEN ? AND ? + GROUP BY + date_trunc('day', created_at) + ORDER BY + date; + "#)?; + let scrobbles = stmt.query_map([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_scrobbles_per_month(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(GetScrobblesPerDayParams::default().start.unwrap()); + let end = params.end.unwrap_or(GetScrobblesPerDayParams::default().end.unwrap()); + let did = params.user_did; + + let conn = conn.lock().unwrap(); + match did { + Some(did) => { + let mut stmt = conn.prepare(r#" + SELECT + EXTRACT(YEAR FROM created_at) || '-' || + LPAD(EXTRACT(MONTH FROM created_at)::VARCHAR, 2, '0') AS year_month, + COUNT(*) AS count + FROM + scrobbles + LEFT JOIN users u ON scrobbles.user_id = u.id + WHERE + u.did = ? + AND created_at BETWEEN ? AND ? + GROUP BY + EXTRACT(YEAR FROM created_at), + EXTRACT(MONTH FROM created_at) + ORDER BY + year_month; + "#)?; + let scrobbles = stmt.query_map([did, start, end], |row| { + Ok(ScrobblesPerMonth { + year_month: row.get(0)?, + count: row.get(1)?, + }) + })?; + let scrobbles: Result, _> = scrobbles.collect(); + Ok(HttpResponse::Ok().json(scrobbles?)) + }, + None => { + let mut stmt = conn.prepare(r#" + SELECT + EXTRACT(YEAR FROM created_at) || '-' || + LPAD(EXTRACT(MONTH FROM created_at)::VARCHAR, 2, '0') AS year_month, + COUNT(*) AS count + FROM + scrobbles + WHERE + created_at BETWEEN ? AND ? + GROUP BY + EXTRACT(YEAR FROM created_at), + EXTRACT(MONTH FROM created_at) + ORDER BY + year_month; + "#)?; + let scrobbles = stmt.query_map([start, end], |row| { + Ok(ScrobblesPerMonth { + year_month: row.get(0)?, + count: row.get(1)?, + }) + })?; + let scrobbles: Result, _> = scrobbles.collect(); + Ok(HttpResponse::Ok().json(scrobbles?)) + } + } +} + +pub async fn get_scrobbles_per_year(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(GetScrobblesPerDayParams::default().start.unwrap()); + let end = params.end.unwrap_or(GetScrobblesPerDayParams::default().end.unwrap()); + let did = params.user_did; + + let conn = conn.lock().unwrap(); + match did { + Some(did) => { + let mut stmt = conn.prepare(r#" + SELECT + EXTRACT(YEAR FROM created_at) AS year, + COUNT(*) AS count + FROM + scrobbles + LEFT JOIN users u ON scrobbles.user_id = u.id + WHERE + u.did = ? + AND created_at BETWEEN ? AND ? + GROUP BY + EXTRACT(YEAR FROM created_at) + ORDER BY + year; + "#)?; + let scrobbles = stmt.query_map([did, start, end], |row| { + Ok(ScrobblesPerYear { + year: row.get(0)?, + count: row.get(1)?, + }) + })?; + let scrobbles: Result, _> = scrobbles.collect(); + Ok(HttpResponse::Ok().json(scrobbles?)) + }, + None => { + let mut stmt = conn.prepare(r#" + SELECT + EXTRACT(YEAR FROM created_at) AS year, + COUNT(*) AS count + FROM + scrobbles + WHERE + created_at BETWEEN ? AND ? + GROUP BY + EXTRACT(YEAR FROM created_at) + ORDER BY + year; + "#)?; + let scrobbles = stmt.query_map([start, end], |row| { + Ok(ScrobblesPerYear { + year: 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 new file mode 100644 index 00000000..2e283525 --- /dev/null +++ b/crates/analytics/src/handlers/tracks.rs @@ -0,0 +1,344 @@ +use std::sync::{Arc, Mutex}; + +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 crate::read_payload; + +pub async fn get_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 did = params.user_did; + + let conn = conn.lock().unwrap(); + match did { + Some(did) => { + let mut stmt = conn.prepare(r#" + SELECT + t.id, + t.title, + t.artist, + t.album_artist, + t.album_art, + t.album, + t.track_number, + t.duration, + t.mb_id, + t.youtube_link, + t.spotify_link, + t.tidal_link, + t.apple_music_link, + t.sha256, + t.composer, + t.genre, + t.disc_number, + t.label, + t.uri, + t.copyright_message, + t.artist_uri, + t.album_uri, + t.created_at, + 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 = ? + ORDER BY t.title ASC + OFFSET ? + LIMIT ?; + "#)?; + let tracks = stmt.query_map([did, 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_art: row.get(4)?, + album: row.get(5)?, + track_number: row.get(6)?, + duration: row.get(7)?, + mb_id: row.get(8)?, + youtube_link: row.get(9)?, + spotify_link: row.get(10)?, + tidal_link: row.get(11)?, + apple_music_link: row.get(12)?, + sha256: row.get(13)?, + composer: row.get(14)?, + genre: row.get(15)?, + disc_number: row.get(16)?, + label: row.get(17)?, + uri: row.get(18)?, + copyright_message: row.get(19)?, + artist_uri: row.get(20)?, + album_uri: row.get(21)?, + created_at: row.get(22)?, + ..Default::default() + }) + })?; + let tracks: Result, _> = tracks.collect(); + Ok(HttpResponse::Ok().json(tracks?)) + }, + None => { + let mut stmt = conn.prepare(r#" + SELECT + id, + title, + artist, + album_artist, + album_art, + album, + track_number, + duration, + mb_id, + youtube_link, + spotify_link, + tidal_link, + apple_music_link, + sha256, + composer, + genre, + disc_number, + label, + uri, + copyright_message, + artist_uri, + album_uri, + created_at, + FROM tracks + ORDER BY title ASC + OFFSET ? + LIMIT ?; + "#)?; + let tracks = stmt.query_map([limit, offset], |row| { + Ok(Track { + id: row.get(0)?, + title: row.get(1)?, + artist: row.get(2)?, + album_artist: row.get(3)?, + album_art: row.get(4)?, + album: row.get(5)?, + track_number: row.get(6)?, + duration: row.get(7)?, + mb_id: row.get(8)?, + youtube_link: row.get(9)?, + spotify_link: row.get(10)?, + tidal_link: row.get(11)?, + apple_music_link: row.get(12)?, + sha256: row.get(13)?, + composer: row.get(14)?, + genre: row.get(15)?, + disc_number: row.get(16)?, + label: row.get(17)?, + uri: row.get(18)?, + copyright_message: row.get(19)?, + artist_uri: row.get(20)?, + album_uri: row.get(21)?, + created_at: row.get(22)?, + ..Default::default() + }) + })?; + let tracks: Result, _> = tracks.collect(); + Ok(HttpResponse::Ok().json(tracks?)) + } + } +} + +pub async fn get_loved_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 did = params.user_did; + + let conn = conn.lock().unwrap(); + let mut stmt = conn.prepare(r#" + SELECT + t.id, + t.title, + t.artist, + t.album, + t.album_artist, + t.album_art, + t.album_uri, + t.artist_uri, + t.composer, + t.copyright_message, + t.disc_number, + t.duration, + t.track_number, + t.label, + t.spotify_link, + t.tidal_link, + t.youtube_link, + t.apple_music_link, + t.sha256, + t.uri, + u.handle, + u.did, + l.created_at + 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 = ? + ORDER BY l.created_at DESC + OFFSET ? + LIMIT ?; + "#)?; + let loved_tracks = stmt.query_map([did, limit.to_string(), offset.to_string()], |row| { + Ok(Track { + id: row.get(0)?, + title: row.get(1)?, + artist: row.get(2)?, + album: row.get(3)?, + album_artist: row.get(4)?, + album_art: row.get(5)?, + album_uri: row.get(6)?, + artist_uri: row.get(7)?, + composer: row.get(8)?, + copyright_message: row.get(9)?, + disc_number: row.get(10)?, + duration: row.get(11)?, + track_number: row.get(12)?, + label: row.get(13)?, + spotify_link: row.get(14)?, + tidal_link: row.get(15)?, + youtube_link: row.get(16)?, + apple_music_link: row.get(17)?, + sha256: row.get(18)?, + uri: row.get(19)?, + handle: row.get(20)?, + did: row.get(21)?, + created_at: row.get(22)?, + ..Default::default() + }) + })?; + let loved_tracks: Result, _> = loved_tracks.collect(); + Ok(HttpResponse::Ok().json(loved_tracks?)) +} + +pub async fn get_top_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 did = params.user_did; + + let conn = conn.lock().unwrap(); + match did { + Some(did) => { + 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.created_at, + COUNT(*) AS play_count, + COUNT(DISTINCT s.user_id) AS unique_listeners + FROM scrobbles s + LEFT JOIN tracks t ON s.track_id = t.id + 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 = ? + 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| { + 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)?, + created_at: row.get(13)?, + play_count: row.get(14)?, + unique_listeners: row.get(15)?, + ..Default::default() + }) + })?; + let top_tracks: Result, _> = top_tracks.collect(); + Ok(HttpResponse::Ok().json(top_tracks?)) + }, + None => { + 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.created_at, + COUNT(*) AS play_count, + COUNT(DISTINCT s.user_id) AS unique_listeners + FROM scrobbles s + LEFT JOIN tracks t ON s.track_id = t.id + LEFT JOIN artists ar ON s.artist_id = ar.id + LEFT JOIN albums a ON s.album_id = a.id + WHERE s.track_id IS NOT NULL AND s.artist_id IS NOT NULL AND s.album_id IS NOT NULL + 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([limit, offset], |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)?, + created_at: row.get(13)?, + play_count: row.get(14)?, + unique_listeners: row.get(15)?, + ..Default::default() + }) + })?; + let top_tracks: Result, _> = top_tracks.collect(); + Ok(HttpResponse::Ok().json(top_tracks?)) + } + } +} \ No newline at end of file diff --git a/crates/analytics/src/main.rs b/crates/analytics/src/main.rs index 7e4c099f..870b5348 100644 --- a/crates/analytics/src/main.rs +++ b/crates/analytics/src/main.rs @@ -1,5 +1,5 @@ use core::create_tables; -use std::env; +use std::{env, sync::{Arc, Mutex}}; use clap::Command; use cmd::{serve::serve, sync::sync}; @@ -10,7 +10,8 @@ use dotenv::dotenv; pub mod types; pub mod xata; pub mod cmd; -pub mod core; +pub mod core; +pub mod handlers; fn cli() -> Command { Command::new("analytics") @@ -37,11 +38,12 @@ async fn main() -> Result<(), Box> { 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?, + Some(("sync", _)) => sync(conn, &pool).await?, + Some(("serve", _)) => serve(conn).await?, + _ => serve(conn).await?, } Ok(()) diff --git a/crates/analytics/src/types/album.rs b/crates/analytics/src/types/album.rs index fe97a74c..ee81033b 100644 --- a/crates/analytics/src/types/album.rs +++ b/crates/analytics/src/types/album.rs @@ -1,32 +1,45 @@ -use actix_web::{body::BoxBody, http::header::ContentType, HttpRequest, HttpResponse, Responder}; use serde::{Deserialize, Serialize}; -#[derive(Debug, Serialize, Deserialize)] +use super::pagination::Pagination; + +#[derive(Debug, Serialize, Deserialize, Default)] pub struct Album { pub id: String, pub title: String, pub artist: String, + #[serde(skip_serializing_if = "Option::is_none")] pub release_date: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub album_art: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub year: 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, + #[serde(skip_serializing_if = "Option::is_none")] pub apple_music_link: Option, pub sha256: String, + #[serde(skip_serializing_if = "Option::is_none")] pub uri: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub artist_uri: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub play_count: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub unique_listeners: Option, } -impl Responder for Album { - type Body = BoxBody; - - fn respond_to(self, _req: &HttpRequest) -> HttpResponse { - let body = serde_json::to_string(&self).unwrap(); +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct GetAlbumsParams { + pub user_did: Option, + pub pagination: Option, +} - // Create response and set content type - HttpResponse::Ok() - .content_type(ContentType::json()) - .body(body) - } +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct GetTopAlbumsParams { + pub user_did: Option, + pub pagination: Option, } diff --git a/crates/analytics/src/types/album_track.rs b/crates/analytics/src/types/album_track.rs deleted file mode 100644 index e69de29b..00000000 diff --git a/crates/analytics/src/types/artist.rs b/crates/analytics/src/types/artist.rs index f4fcd2bd..12dee892 100644 --- a/crates/analytics/src/types/artist.rs +++ b/crates/analytics/src/types/artist.rs @@ -1,33 +1,46 @@ -use actix_web::{body::BoxBody, http::header::ContentType, HttpRequest, HttpResponse, Responder}; +use super::pagination::Pagination; use chrono::NaiveDate; use serde::{Deserialize, Serialize}; -#[derive(Debug, Serialize, Deserialize)] +#[derive(Debug, Serialize, Deserialize, Default)] pub struct Artist { pub id: String, pub name: String, + #[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 picture: Option, pub sha256: String, + #[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, + #[serde(skip_serializing_if = "Option::is_none")] pub apple_music_link: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub uri: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub play_count: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub unique_listeners: Option, } -impl Responder for Artist { - type Body = BoxBody; - - fn respond_to(self, _req: &HttpRequest) -> HttpResponse { - let body = serde_json::to_string(&self).unwrap(); +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct GetArtistsParams { + pub user_did: Option, + pub pagination: Option, +} - // Create response and set content type - HttpResponse::Ok() - .content_type(ContentType::json()) - .body(body) - } +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct GetTopArtistsParams { + pub user_did: Option, + pub pagination: Option, } diff --git a/crates/analytics/src/types/artist_album.rs b/crates/analytics/src/types/artist_album.rs deleted file mode 100644 index e69de29b..00000000 diff --git a/crates/analytics/src/types/artist_track.rs b/crates/analytics/src/types/artist_track.rs deleted file mode 100644 index e69de29b..00000000 diff --git a/crates/analytics/src/types/filters.rs b/crates/analytics/src/types/filters.rs new file mode 100644 index 00000000..13877650 --- /dev/null +++ b/crates/analytics/src/types/filters.rs @@ -0,0 +1,8 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Serialize, Deserialize, Default)] +pub struct Filters { + pub user_did: Option, + pub order_by: Option, + pub asc: Option, +} diff --git a/crates/analytics/src/types/loved_track.rs b/crates/analytics/src/types/loved_track.rs deleted file mode 100644 index e69de29b..00000000 diff --git a/crates/analytics/src/types/mod.rs b/crates/analytics/src/types/mod.rs index c7f2208b..9da51358 100644 --- a/crates/analytics/src/types/mod.rs +++ b/crates/analytics/src/types/mod.rs @@ -1,12 +1,9 @@ pub mod album; -pub mod album_track; pub mod artist; -pub mod artist_track; +pub mod filters; +pub mod pagination; pub mod playlist; -pub mod playlist_track; pub mod scrobble; +pub mod stats; pub mod track; pub mod user; -pub mod user_album; -pub mod user_artist; -pub mod user_playlist; diff --git a/crates/analytics/src/types/pagination.rs b/crates/analytics/src/types/pagination.rs new file mode 100644 index 00000000..7adc249f --- /dev/null +++ b/crates/analytics/src/types/pagination.rs @@ -0,0 +1,16 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Pagination { + pub skip: Option, + pub take: Option, +} + +impl Default for Pagination { + fn default() -> Self { + Pagination { + skip: Some(0), + take: Some(20), + } + } +} diff --git a/crates/analytics/src/types/playlist.rs b/crates/analytics/src/types/playlist.rs index ecf44d99..91ea1bb9 100644 --- a/crates/analytics/src/types/playlist.rs +++ b/crates/analytics/src/types/playlist.rs @@ -1,8 +1,7 @@ -use actix_web::{body::BoxBody, http::header::ContentType, HttpRequest, HttpResponse, Responder}; use chrono::NaiveDateTime; use serde::{Deserialize, Serialize}; -#[derive(Debug, Serialize, Deserialize)] +#[derive(Debug, Serialize, Deserialize, Default)] pub struct Playlist { pub id: String, pub name: String, @@ -12,17 +11,5 @@ pub struct Playlist { pub updated_at: NaiveDateTime, pub uri: Option, pub created_by: String, -} - -impl Responder for Playlist { - type Body = BoxBody; - - fn respond_to(self, _req: &HttpRequest) -> HttpResponse { - let body = serde_json::to_string(&self).unwrap(); - - // Create response and set content type - HttpResponse::Ok() - .content_type(ContentType::json()) - .body(body) - } + pub listeners: Option, } diff --git a/crates/analytics/src/types/playlist_track.rs b/crates/analytics/src/types/playlist_track.rs deleted file mode 100644 index e69de29b..00000000 diff --git a/crates/analytics/src/types/scrobble.rs b/crates/analytics/src/types/scrobble.rs index 99d92eb3..5fb04c28 100644 --- a/crates/analytics/src/types/scrobble.rs +++ b/crates/analytics/src/types/scrobble.rs @@ -1,7 +1,9 @@ -use chrono::NaiveDateTime; +use chrono::{NaiveDate, NaiveDateTime}; use serde::{Deserialize, Serialize}; -#[derive(Debug, Serialize, Deserialize)] +use super::pagination::Pagination; + +#[derive(Debug, Serialize, Deserialize, Default)] pub struct Scrobble { pub id: String, pub user_id: String, @@ -11,3 +13,44 @@ pub struct Scrobble { pub uri: Option, pub created_at: NaiveDateTime, } + +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct ScrobbleTrack { + pub id: String, + pub track_id: String, + pub title: String, + pub artist: String, + pub album_artist: String, + pub album_art: Option, + pub album: String, + pub handle: String, + pub uri: Option, + pub track_uri: Option, + pub artist_uri: Option, + pub album_uri: Option, + pub created_at: NaiveDateTime, +} + +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct ScrobblesPerDay { + pub date: NaiveDate, + pub count: i32, +} + +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct ScrobblesPerMonth { + pub year_month: String, + pub count: i32, +} + +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct ScrobblesPerYear { + pub year: i32, + pub count: i32, +} + +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct GetScrobblesParams { + pub user_did: Option, + pub pagination: Option, +} diff --git a/crates/analytics/src/types/stats.rs b/crates/analytics/src/types/stats.rs new file mode 100644 index 00000000..bd636f0e --- /dev/null +++ b/crates/analytics/src/types/stats.rs @@ -0,0 +1,70 @@ +use super::pagination::Pagination; + +use chrono::{Datelike, Duration, NaiveDate, Utc}; +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct GetStatsParams { + pub user_did: String, + pub pagination: Option, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct GetScrobblesPerDayParams { + pub user_did: Option, + pub start: Option, + pub end: Option, +} + +impl Default for GetScrobblesPerDayParams { + fn default() -> Self { + let current_date = Utc::now().naive_utc(); + let date_30_days_ago = current_date - Duration::days(30); + + GetScrobblesPerDayParams { + user_did: None, + start: Some(date_30_days_ago.to_string()), + end: Some(current_date.to_string()), + } + } +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct GetScrobblesPerMonthParams { + pub user_did: Option, + pub start: Option, + pub end: Option, +} + +impl Default for GetScrobblesPerMonthParams { + fn default() -> Self { + let current_date = Utc::now().naive_utc(); + let january = NaiveDate::from_ymd_opt(current_date.year(), 1, 1).unwrap(); + + GetScrobblesPerMonthParams { + user_did: None, + start: Some(january.to_string()), + end: Some(current_date.to_string()), + } + } +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct GetScrobblesPerYearParams { + pub user_did: Option, + pub start: Option, + pub end: Option, +} + +impl Default for GetScrobblesPerYearParams { + fn default() -> Self { + let current_date = Utc::now().naive_utc(); + let start = NaiveDate::from_ymd_opt(2025, 1, 1).unwrap(); + + GetScrobblesPerYearParams { + user_did: None, + start: Some(start.to_string()), + end: Some(current_date.to_string()), + } + } +} diff --git a/crates/analytics/src/types/track.rs b/crates/analytics/src/types/track.rs index 1dd8aebc..58e29e0f 100644 --- a/crates/analytics/src/types/track.rs +++ b/crates/analytics/src/types/track.rs @@ -1,44 +1,72 @@ -use actix_web::{body::BoxBody, http::header::ContentType, HttpRequest, HttpResponse, Responder}; +use super::pagination::Pagination; + use chrono::NaiveDateTime; use serde::{Deserialize, Serialize}; -#[derive(Debug, Serialize, Deserialize)] +#[derive(Debug, Serialize, Deserialize, Default)] pub struct Track { pub id: String, pub title: String, pub artist: String, pub album_artist: String, + #[serde(skip_serializing_if = "Option::is_none")] pub album_art: Option, pub album: String, pub track_number: i32, pub duration: i32, + #[serde(skip_serializing_if = "Option::is_none")] pub mb_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub youtube_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 apple_music_link: Option, pub sha256: String, + #[serde(skip_serializing_if = "Option::is_none")] pub lyrics: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub composer: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub genre: Option, pub disc_number: i32, + #[serde(skip_serializing_if = "Option::is_none")] pub copyright_message: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub label: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub uri: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub artist_uri: Option, + #[serde(skip_serializing_if = "Option::is_none")] pub album_uri: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub handle: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub did: Option, pub created_at: NaiveDateTime, + #[serde(skip_serializing_if = "Option::is_none")] + pub play_count: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub unique_listeners: Option, } -impl Responder for Track { - type Body = BoxBody; +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct GetTracksParams { + pub user_did: Option, + pub pagination: Option, +} - fn respond_to(self, _req: &HttpRequest) -> HttpResponse { - let body = serde_json::to_string(&self).unwrap(); +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct GetTopTracksParams { + pub user_did: Option, + pub pagination: Option, +} - // Create response and set content type - HttpResponse::Ok() - .content_type(ContentType::json()) - .body(body) - } +#[derive(Debug, Serialize, Deserialize, Default)] +pub struct GetLovedTracksParams { + pub user_did: String, + pub pagination: Option, } diff --git a/crates/analytics/src/types/user.rs b/crates/analytics/src/types/user.rs index 6c5cc786..b77911c7 100644 --- a/crates/analytics/src/types/user.rs +++ b/crates/analytics/src/types/user.rs @@ -1,7 +1,6 @@ -use actix_web::{body::BoxBody, http::header::ContentType, HttpRequest, HttpResponse, Responder}; use serde::{Deserialize, Serialize}; -#[derive(Debug, Serialize, Deserialize)] +#[derive(Debug, Serialize, Deserialize, Default)] pub struct User { pub id: String, pub display_name: String, @@ -9,16 +8,3 @@ pub struct User { pub handle: String, pub avatar: String, } - -impl Responder for User { - type Body = BoxBody; - - fn respond_to(self, _req: &HttpRequest) -> HttpResponse { - let body = serde_json::to_string(&self).unwrap(); - - // Create response and set content type - HttpResponse::Ok() - .content_type(ContentType::json()) - .body(body) - } -} diff --git a/crates/analytics/src/types/user_album.rs b/crates/analytics/src/types/user_album.rs deleted file mode 100644 index e69de29b..00000000 diff --git a/crates/analytics/src/types/user_artist.rs b/crates/analytics/src/types/user_artist.rs deleted file mode 100644 index e69de29b..00000000 diff --git a/crates/analytics/src/types/user_playlist.rs b/crates/analytics/src/types/user_playlist.rs deleted file mode 100644 index e69de29b..00000000 diff --git a/crates/analytics/src/types/user_track.rs b/crates/analytics/src/types/user_track.rs deleted file mode 100644 index e69de29b..00000000 -- 2.51.2