diff --git a/crates/analytics/src/core.rs b/crates/analytics/src/core.rs index a3f6d6a5..951a54cb 100644 --- a/crates/analytics/src/core.rs +++ b/crates/analytics/src/core.rs @@ -1,4 +1,7 @@ -use std::sync::{Arc, Mutex}; +use std::{ + env, + sync::{Arc, Mutex}, +}; use anyhow::Error; use duckdb::{params, Connection}; @@ -180,6 +183,88 @@ pub async fn create_tables(conn: &Connection) -> Result<(), Error> { ", )?; + match conn.execute("ALTER TABLE artists ADD COLUMN genres VARCHAR[]", []) { + Ok(_) => tracing::info!("Added genres column to artists table"), + Err(e) => tracing::warn!("Could not add genres column to artists table: {}", e), + } + + Ok(()) +} + +pub async fn update_artist_genres( + conn: Arc>, + pool: &Pool, +) -> Result<(), Error> { + if env::var("UPDATE_ARTIST_GENRES").is_err() { + tracing::info!("Skipping update_artist_genres as UPDATE_ARTIST_GENRES is not set"); + return Ok(()); + } + + let artists: Vec = sqlx::query_as( + r#" + SELECT * FROM artists + "#, + ) + .fetch_all(pool) + .await?; + + let conn = conn.lock().unwrap(); + + conn.execute_batch( + " + CREATE TABLE user_artists_new AS SELECT * FROM user_artists; + CREATE TABLE artist_tracks_new AS SELECT * FROM artist_tracks; + CREATE TABLE artist_albums_new AS SELECT * FROM artist_albums; + CREATE TABLE scrobbles_new AS SELECT * FROM scrobbles; + CREATE TABLE loved_tracks_new AS SELECT * FROM loved_tracks; + DROP TABLE user_artists; + DROP TABLE artist_tracks; + DROP TABLE artist_albums; + DROP TABLE scrobbles; + DROP TABLE loved_tracks; + ALTER TABLE user_artists_new RENAME TO user_artists; + ALTER TABLE artist_tracks_new RENAME TO artist_tracks; + ALTER TABLE artist_albums_new RENAME TO artist_albums; + ALTER TABLE scrobbles_new RENAME TO scrobbles; + ALTER TABLE loved_tracks_new RENAME TO loved_tracks; + ", + )?; + + for (i, artist) in artists.clone().into_iter().enumerate() { + let genres_array = artist + .genres + .as_ref() + .map(|tags| { + tags.iter() + .map(|tag| format!("'{}'", tag.replace("'", "''"))) + .collect::>() + .join(", ") + }) + .unwrap_or_default() + .trim() + .to_string(); + + if genres_array.is_empty() { + continue; + } + + tracing::info!(artist = i, name = %artist.name.bright_green(), genres = %genres_array, "Updating artist genres"); + + match conn.execute( + &format!( + "UPDATE artists SET genres = [{}] WHERE id = ? AND genres IS NULL", + genres_array + ), + params![artist.xata_id], + ) { + Ok(_) => (), + Err(e) => { + tracing::error!(error = %e, genres = %genres_array, "Error updating artist >> ") + } + } + } + + tracing::info!(artists = artists.len(), "Updated artist genres"); Ok(()) } diff --git a/crates/analytics/src/handlers/artists.rs b/crates/analytics/src/handlers/artists.rs index 0be49544..e3f40f43 100644 --- a/crates/analytics/src/handlers/artists.rs +++ b/crates/analytics/src/handlers/artists.rs @@ -10,7 +10,7 @@ use crate::types::{ }; use actix_web::{web, HttpRequest, HttpResponse}; use anyhow::Error; -use duckdb::Connection; +use duckdb::{types::Value, Connection}; use tokio_stream::StreamExt; use crate::read_payload; @@ -59,6 +59,7 @@ pub async fn get_artists( let artists = stmt.query_map( [&did, &did, &limit.to_string(), &offset.to_string()], |row| { + let genres = extract_genres_from_value(row.get(13)?); Ok(Artist { id: row.get(0)?, name: row.get(1)?, @@ -73,8 +74,9 @@ pub async fn get_artists( youtube_link: row.get(10)?, apple_music_link: row.get(11)?, uri: row.get(12)?, - play_count: row.get(13)?, - unique_listeners: row.get(14)?, + genres, + play_count: row.get(14)?, + unique_listeners: row.get(15)?, }) }, )?; @@ -84,6 +86,7 @@ pub async fn get_artists( } None => { let artists = stmt.query_map([limit, offset], |row| { + let genres = extract_genres_from_value(row.get(13)?); Ok(Artist { id: row.get(0)?, name: row.get(1)?, @@ -98,8 +101,9 @@ pub async fn get_artists( youtube_link: row.get(10)?, apple_music_link: row.get(11)?, uri: row.get(12)?, - play_count: row.get(13)?, - unique_listeners: row.get(14)?, + genres, + play_count: row.get(14)?, + unique_listeners: row.get(15)?, }) })?; @@ -131,6 +135,7 @@ pub async fn get_top_artists( ar.picture AS picture, ar.sha256 AS sha256, ar.uri AS uri, + ar.genres AS genres, COUNT(*) AS play_count, COUNT(DISTINCT s.user_id) AS unique_listeners FROM @@ -142,7 +147,7 @@ pub async fn get_top_artists( WHERE 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 + s.artist_id, ar.name, ar.uri, ar.picture, ar.sha256, ar.genres ORDER BY play_count DESC OFFSET ? @@ -157,6 +162,7 @@ pub async fn get_top_artists( ar.picture AS picture, ar.sha256 AS sha256, ar.uri AS uri, + ar.genres AS genres, COUNT(*) AS play_count, COUNT(DISTINCT s.user_id) AS unique_listeners FROM @@ -166,7 +172,7 @@ pub async fn get_top_artists( WHERE s.artist_id IS NOT NULL GROUP BY - s.artist_id, ar.name, ar.uri, ar.picture, ar.sha256 + s.artist_id, ar.name, ar.uri, ar.picture, ar.sha256, ar.genres ORDER BY play_count DESC OFFSET ? @@ -180,6 +186,7 @@ pub async fn get_top_artists( let artists = stmt.query_map( [&did, &did, &limit.to_string(), &offset.to_string()], |row| { + let genres = extract_genres_from_value(row.get(5)?); Ok(Artist { id: row.get(0)?, name: row.get(1)?, @@ -194,8 +201,9 @@ pub async fn get_top_artists( youtube_link: None, apple_music_link: None, uri: row.get(4)?, - play_count: Some(row.get(5)?), - unique_listeners: Some(row.get(6)?), + genres, + play_count: Some(row.get(6)?), + unique_listeners: Some(row.get(7)?), }) }, )?; @@ -205,6 +213,7 @@ pub async fn get_top_artists( } None => { let artists = stmt.query_map([limit, offset], |row| { + let genres = extract_genres_from_value(row.get(5)?); Ok(Artist { id: row.get(0)?, name: row.get(1)?, @@ -219,8 +228,9 @@ pub async fn get_top_artists( youtube_link: None, apple_music_link: None, uri: row.get(4)?, - play_count: Some(row.get(5)?), - unique_listeners: Some(row.get(6)?), + genres, + play_count: Some(row.get(6)?), + unique_listeners: Some(row.get(7)?), }) })?; @@ -491,3 +501,20 @@ pub async fn get_artist_listeners( let listeners: Result, _> = listeners.collect(); Ok(HttpResponse::Ok().json(listeners?)) } + +fn extract_genres_from_value(value: Value) -> Option> { + match value { + Value::Null => None, + Value::List(items) => { + let genres: Vec = items + .into_iter() + .filter_map(|item| match item { + Value::Text(s) => Some(s), + _ => None, + }) + .collect(); + Some(genres) + } + _ => None, + } +} diff --git a/crates/analytics/src/lib.rs b/crates/analytics/src/lib.rs index 5179fb4f..149956a4 100644 --- a/crates/analytics/src/lib.rs +++ b/crates/analytics/src/lib.rs @@ -9,7 +9,7 @@ use anyhow::Error; use duckdb::Connection; use sqlx::postgres::PgPoolOptions; -use crate::core::create_tables; +use crate::core::{create_tables, update_artist_genres}; pub mod cmd; pub mod core; @@ -23,7 +23,14 @@ pub async fn serve() -> Result<(), Error> { create_tables(&conn).await?; + let pool = PgPoolOptions::new() + .max_connections(5) + .connect(&env::var("XATA_POSTGRES_URL")?) + .await?; + let conn = Arc::new(Mutex::new(conn)); + update_artist_genres(conn.clone(), &pool).await?; + export_parquets(conn.clone()); cmd::serve::serve(conn).await?; diff --git a/crates/analytics/src/subscriber/mod.rs b/crates/analytics/src/subscriber/mod.rs index a64fe50e..abd75866 100644 --- a/crates/analytics/src/subscriber/mod.rs +++ b/crates/analytics/src/subscriber/mod.rs @@ -207,7 +207,8 @@ pub async fn save_scrobble( let conn = conn.lock().unwrap(); match conn.execute( - "INSERT INTO artists ( + &format!( + "INSERT INTO artists ( id, name, biography, @@ -220,7 +221,8 @@ pub async fn save_scrobble( tidal_link, youtube_link, apple_music_link, - uri + uri, + [{}] ) VALUES ( ?, ?, @@ -236,6 +238,18 @@ pub async fn save_scrobble( ?, ? )", + payload + .scrobble + .artist_id + .genres + .as_ref() + .map(|genres| genres + .iter() + .map(|g| format!("'{}'", g)) + .collect::>() + .join(", ")) + .unwrap_or_default() + ), params![ payload.scrobble.artist_id.xata_id, payload.scrobble.artist_id.name, diff --git a/crates/analytics/src/subscriber/types.rs b/crates/analytics/src/subscriber/types.rs index bf3ca9de..90b695fa 100644 --- a/crates/analytics/src/subscriber/types.rs +++ b/crates/analytics/src/subscriber/types.rs @@ -151,6 +151,8 @@ pub struct ArtistId { pub tidal_link: Option, #[serde(skip_serializing_if = "Option::is_none")] pub youtube_link: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub genres: Option>, } #[derive(Debug, Serialize, Deserialize, Clone)] diff --git a/crates/analytics/src/types/artist.rs b/crates/analytics/src/types/artist.rs index 1e3e8c24..d38e63f3 100644 --- a/crates/analytics/src/types/artist.rs +++ b/crates/analytics/src/types/artist.rs @@ -31,6 +31,8 @@ pub struct Artist { pub play_count: Option, #[serde(skip_serializing_if = "Option::is_none")] pub unique_listeners: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub genres: Option>, } #[derive(Debug, Serialize, Deserialize, Default)] diff --git a/crates/analytics/src/xata/artist.rs b/crates/analytics/src/xata/artist.rs index 25187ec5..f85f8410 100644 --- a/crates/analytics/src/xata/artist.rs +++ b/crates/analytics/src/xata/artist.rs @@ -18,6 +18,27 @@ pub struct Artist { pub youtube_link: Option, pub apple_music_link: Option, pub uri: Option, + pub genres: Option>, #[serde(with = "chrono::serde::ts_seconds")] pub xata_createdat: DateTime, } + +#[derive(Debug, sqlx::FromRow, Deserialize, Clone)] +pub struct ArtistWithoutDate { + pub xata_id: String, + pub name: String, + pub biography: Option, + #[serde(with = "chrono::serde::ts_seconds_option")] + pub born: Option>, + pub born_in: Option, + #[serde(with = "chrono::serde::ts_seconds_option")] + pub died: Option>, + pub picture: Option, + pub sha256: String, + pub spotify_link: Option, + pub tidal_link: Option, + pub youtube_link: Option, + pub apple_music_link: Option, + pub uri: Option, + pub genres: Option>, +}