From e74e74ee807d39eafea3a75e4cce512d7ccc3b03 Mon Sep 17 00:00:00 2001 From: Tsiry Sandratraina Date: Mon, 3 Mar 2025 00:26:30 +0300 Subject: [PATCH] rocksky: publish new track to nats --- crates/analytics/src/subscriber/mod.rs | 216 +++++++++++++++++- crates/analytics/src/subscriber/types.rs | 40 ++++ crates/analytics/src/xata/track.rs | 4 +- .../rocksky-auth/src/tracks/tracks.service.ts | 27 ++- 4 files changed, 277 insertions(+), 10 deletions(-) diff --git a/crates/analytics/src/subscriber/mod.rs b/crates/analytics/src/subscriber/mod.rs index 483a5440..9e01717a 100644 --- a/crates/analytics/src/subscriber/mod.rs +++ b/crates/analytics/src/subscriber/mod.rs @@ -4,7 +4,7 @@ use async_nats::{connect, Client}; use duckdb::{params, Connection}; use owo_colors::OwoColorize; use tokio_stream::StreamExt; -use types::{LikePayload, ScrobblePayload, UnlikePayload}; +use types::{LikePayload, NewTrackPayload, ScrobblePayload, UnlikePayload}; pub mod types; @@ -54,6 +54,38 @@ pub fn on_scrobble(nc: Arc>, conn: Arc>) { }); } +pub fn on_new_track(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.track".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_track(conn.clone(), payload.clone()).await { + Ok(_) => println!("Song saved successfully for {}", payload.track.title.cyan()), + Err(e) => eprintln!("Error saving song: {}", e), + } + }, + Err(e) => { + eprintln!("Error parsing payload: {}", e); + } + } + } + + Ok::<(), Error>(()) + })?; + + Ok::<(), Error>(()) + }); +} + pub fn on_like(nc: Arc>, conn: Arc>) { thread::spawn(move || { @@ -448,6 +480,188 @@ pub async fn save_scrobble(conn: Arc>, payload: ScrobblePayloa Ok(()) } +pub async fn save_track(conn: Arc>, payload: NewTrackPayload) -> Result<(), Error> { + let conn = conn.lock().unwrap(); + + 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.track.xata_id, + payload.track.title, + payload.track.artist, + payload.track.album_artist, + payload.track.album_art, + payload.track.album, + payload.track.track_number, + payload.track.duration, + payload.track.mb_id, + payload.track.youtube_link, + payload.track.spotify_link, + payload.track.tidal_link, + payload.track.apple_music_link, + payload.track.sha256, + payload.track.lyrics, + payload.track.composer, + payload.track.genre, + payload.track.disc_number, + payload.track.copyright_message, + payload.track.label, + payload.track.uri, + payload.track.artist_uri, + payload.track.album_uri, + payload.track.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 artist_albums (id, artist_id, album_id, created_at) VALUES (?, ?, ?, ?)", + params![ + payload.artist_album.xata_id, + payload.artist_album.artist_id.xata_id, + payload.artist_album.album_id.xata_id, + payload.artist_album.xata_createdat, + ], +) { + Ok(_) => (), + Err(e) => { + if !e.to_string().contains("violates primary key constraint") { + println!("[artist_albums] 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()); + } + } + } + + Ok(()) +} + pub async fn like(conn: Arc>, payload: LikePayload) -> Result<(), Error> { let conn = conn.lock().unwrap(); match conn.execute( diff --git a/crates/analytics/src/subscriber/types.rs b/crates/analytics/src/subscriber/types.rs index d5add7f0..03f49e29 100644 --- a/crates/analytics/src/subscriber/types.rs +++ b/crates/analytics/src/subscriber/types.rs @@ -25,6 +25,17 @@ pub struct UnlikePayload { pub xata_version: i32, } +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct NewTrackPayload { + pub track: Track, + pub user_album: UserAlbum, + pub user_artist: UserArtist, + pub user_track: UserTrack, + pub album_track: AlbumTrack, + pub artist_track: ArtistTrack, + pub artist_album: ArtistAlbum, +} + #[derive(Debug, Serialize, Deserialize, Clone)] pub struct ScrobblePayload { pub scrobble: Scrobble, @@ -36,6 +47,35 @@ pub struct ScrobblePayload { pub artist_album: ArtistAlbum, } +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct Track { + pub xata_id: String, + pub title: String, + pub artist: String, + pub album_artist: String, + pub album_art: Option, + pub album: String, + pub track_number: i32, + pub duration: i32, + pub mb_id: Option, + pub youtube_link: Option, + pub spotify_link: Option, + pub tidal_link: Option, + pub apple_music_link: Option, + pub sha256: String, + pub lyrics: Option, + pub composer: Option, + pub genre: Option, + pub disc_number: i32, + pub copyright_message: Option, + pub label: Option, + pub uri: Option, + pub artist_uri: Option, + pub album_uri: Option, + #[serde(with = "chrono::serde::ts_seconds")] + pub xata_createdat: DateTime, +} + #[derive(Debug, Serialize, Deserialize, Clone)] pub struct Scrobble { pub album_id: AlbumId, diff --git a/crates/analytics/src/xata/track.rs b/crates/analytics/src/xata/track.rs index cd36b095..cc2b0015 100644 --- a/crates/analytics/src/xata/track.rs +++ b/crates/analytics/src/xata/track.rs @@ -1,7 +1,7 @@ use chrono::{DateTime, Utc}; -use serde::Deserialize; +use serde::{Deserialize, Serialize}; -#[derive(Debug, sqlx::FromRow, Deserialize, Clone)] +#[derive(Debug, sqlx::FromRow, Serialize, Deserialize, Clone)] pub struct Track { pub xata_id: String, pub title: String, diff --git a/rockskyapi/rocksky-auth/src/tracks/tracks.service.ts b/rockskyapi/rocksky-auth/src/tracks/tracks.service.ts index 131e0877..a23ce69b 100644 --- a/rockskyapi/rocksky-auth/src/tracks/tracks.service.ts +++ b/rockskyapi/rocksky-auth/src/tracks/tracks.service.ts @@ -28,7 +28,7 @@ export async function saveTrack(ctx: Context, track: Track, agent: Agent) { trackUri = await putSongRecord(track, agent); } - const { xata_id: track_id } = await ctx.client.db.tracks.createOrUpdate( + const newTrack = await ctx.client.db.tracks.createOrUpdate( existingTrack?.xata_id, { title: track.title, @@ -56,6 +56,7 @@ export async function saveTrack(ctx: Context, track: Track, agent: Agent) { spotify_link: track.spotifyLink ? track.spotifyLink : undefined, } ); + const track_id = newTrack.xata_id; const existingArtist = await ctx.client.db.artists .filter( @@ -122,17 +123,20 @@ export async function saveTrack(ctx: Context, track: Track, agent: Agent) { .filter("track_id", equals(track_id)) .getFirst(); - await ctx.client.db.album_tracks.createOrUpdate(existingAlbumTrack?.xata_id, { - album_id, - track_id, - }); + const album_track = await ctx.client.db.album_tracks.createOrUpdate( + existingAlbumTrack?.xata_id, + { + album_id, + track_id, + } + ); const existingArtistTrack = await ctx.client.db.artist_tracks .filter("artist_id", equals(artist_id)) .filter("track_id", equals(track_id)) .getFirst(); - await ctx.client.db.artist_tracks.createOrUpdate( + const artist_track = await ctx.client.db.artist_tracks.createOrUpdate( existingArtistTrack?.xata_id, { artist_id, @@ -145,11 +149,20 @@ export async function saveTrack(ctx: Context, track: Track, agent: Agent) { .filter("album_id", equals(album_id)) .getFirst(); - await ctx.client.db.artist_albums.createOrUpdate( + const artist_album = await ctx.client.db.artist_albums.createOrUpdate( existingArtistAlbum?.xata_id, { artist_id, album_id, } ); + + const message = JSON.stringify({ + track, + album_track, + artist_track, + artist_album, + }); + + ctx.nc.publish("rocksky.track", Buffer.from(message)); } -- 2.51.2