From 6ee0db10116e1b3652ae273e843abca9946aea66 Mon Sep 17 00:00:00 2001 From: Tsiry Sandratraina Date: Sun, 11 Jan 2026 18:23:37 +0300 Subject: [PATCH] Publish scrobble on save and sync URIs Publish the newest scrobble ID after saving so subscribers can react. Refactor JS subscriber to call updateUris(ctx, did) and refreshScrobbles, use the tables namespace for DB queries, add hashing/logging, and limit syncs to 5 recent scrobbles. --- apps/api/src/scripts/sync.ts | 2 +- apps/api/src/subscribers/scrobble.ts | 130 +++++++++------------------ crates/jetstream/src/repo.rs | 17 +++- 3 files changed, 57 insertions(+), 92 deletions(-) diff --git a/apps/api/src/scripts/sync.ts b/apps/api/src/scripts/sync.ts index 7b04a5d7..ac3544a0 100644 --- a/apps/api/src/scripts/sync.ts +++ b/apps/api/src/scripts/sync.ts @@ -12,7 +12,7 @@ import users from "../schema/users"; const args = process.argv.slice(2); -async function updateUris(did: string) { +export async function updateUris(did: string) { // Get scrobbles with track and user data const records = await ctx.db .select({ diff --git a/apps/api/src/subscribers/scrobble.ts b/apps/api/src/subscribers/scrobble.ts index f479f13a..73f5034f 100644 --- a/apps/api/src/subscribers/scrobble.ts +++ b/apps/api/src/subscribers/scrobble.ts @@ -1,9 +1,12 @@ import { consola } from "consola"; import type { Context } from "context"; -import { eq } from "drizzle-orm"; +import { eq, or, desc } from "drizzle-orm"; import _ from "lodash"; import { StringCodec } from "nats"; +import { createHash } from "node:crypto"; import tables from "schema"; +import chalk from "chalk"; +import { publishScrobble } from "nowplaying/nowplaying.service"; export function onNewScrobble(ctx: Context) { const sc = StringCodec(); @@ -15,43 +18,34 @@ export function onNewScrobble(ctx: Context) { .select() .from(tables.scrobbles) .where(eq(tables.scrobbles.id, scrobbleId)) + .leftJoin(tables.users, eq(tables.scrobbles.userId, tables.users.id)) .execute() .then((rows) => rows[0]); if (!result) { consola.info(`Scrobble with ID ${scrobbleId} not found, skipping`); + continue; } + + await updateUris(ctx, result.users.did); + await refreshScrobbles(ctx, result.users.did); } })(); } -/* -import chalk from "chalk"; -import { ctx } from "context"; -import { desc, eq, or } from "drizzle-orm"; -import { createHash } from "node:crypto"; -import { publishScrobble } from "nowplaying/nowplaying.service"; -import albums from "../schema/albums"; -import artists from "../schema/artists"; -import scrobbles from "../schema/scrobbles"; -import tracks from "../schema/tracks"; -import users from "../schema/users"; - -const args = process.argv.slice(2); - -async function updateUris(did: string) { +async function updateUris(ctx: Context, did: string) { // Get scrobbles with track and user data const records = await ctx.db .select({ - track: tracks, - user: users, + track: tables.tracks, + user: tables.users, }) - .from(scrobbles) - .innerJoin(tracks, eq(scrobbles.trackId, tracks.id)) - .innerJoin(users, eq(scrobbles.userId, users.id)) - .where(or(eq(users.did, did), eq(users.handle, did))) - .orderBy(desc(scrobbles.createdAt)) - .limit(process.env.SYNC_SIZE ? parseInt(process.env.SYNC_SIZE, 10) : 20); + .from(tables.scrobbles) + .innerJoin(tables.tracks, eq(tables.scrobbles.trackId, tables.tracks.id)) + .innerJoin(tables.users, eq(tables.scrobbles.userId, tables.users.id)) + .where(or(eq(tables.users.did, did), eq(tables.users.handle, did))) + .orderBy(desc(tables.scrobbles.createdAt)) + .limit(5); for (const { track } of records) { const trackHash = createHash("sha256") @@ -60,8 +54,8 @@ async function updateUris(did: string) { const existingTrack = await ctx.db .select() - .from(tracks) - .where(eq(tracks.sha256, trackHash)) + .from(tables.tracks) + .where(eq(tables.tracks.sha256, trackHash)) .limit(1) .then((rows) => rows[0]); @@ -74,16 +68,16 @@ async function updateUris(did: string) { const album = await ctx.db .select() - .from(albums) - .where(eq(albums.sha256, albumHash)) + .from(tables.albums) + .where(eq(tables.albums.sha256, albumHash)) .limit(1) .then((rows) => rows[0]); if (album) { await ctx.db - .update(tracks) + .update(tables.tracks) .set({ albumUri: album.uri }) - .where(eq(tracks.id, existingTrack.id)); + .where(eq(tables.tracks.id, existingTrack.id)); } } @@ -96,16 +90,16 @@ async function updateUris(did: string) { const artist = await ctx.db .select() - .from(artists) - .where(eq(artists.sha256, artistHash)) + .from(tables.artists) + .where(eq(tables.artists.sha256, artistHash)) .limit(1) .then((rows) => rows[0]); if (artist) { await ctx.db - .update(tracks) + .update(tables.tracks) .set({ artistUri: artist.uri }) - .where(eq(tracks.id, existingTrack.id)); + .where(eq(tables.tracks.id, existingTrack.id)); } } @@ -115,8 +109,8 @@ async function updateUris(did: string) { const album = await ctx.db .select() - .from(albums) - .where(eq(albums.sha256, albumHash)) + .from(tables.albums) + .where(eq(tables.albums.sha256, albumHash)) .limit(1) .then((rows) => rows[0]); @@ -129,69 +123,31 @@ async function updateUris(did: string) { const artist = await ctx.db .select() - .from(artists) - .where(eq(artists.sha256, artistHash)) + .from(tables.artists) + .where(eq(tables.artists.sha256, artistHash)) .limit(1) .then((rows) => rows[0]); if (artist) { await ctx.db - .update(albums) + .update(tables.albums) .set({ artistUri: artist.uri }) - .where(eq(albums.id, album.id)); + .where(eq(tables.albums.id, album.id)); } } } } -if (args.includes("--background")) { - consola.info("Wait for new scrobbles to sync ..."); - const sub = ctx.nc.subscribe("rocksky.user.scrobble.sync"); - for await (const m of sub) { - const did = new TextDecoder().decode(m.data); - // wait for 15 seconds to ensure the scrobble is fully created - await new Promise((resolve) => setTimeout(resolve, 15000)); - consola.info(`Syncing scrobbles ${chalk.magenta(did)} ...`); - await updateUris(did); - - const records = await ctx.db - .select({ - scrobble: scrobbles, - }) - .from(scrobbles) - .innerJoin(users, eq(scrobbles.userId, users.id)) - .where(or(eq(users.did, did), eq(users.handle, did))) - .orderBy(desc(scrobbles.createdAt)) - .limit(5); - - for (const { scrobble } of records) { - consola.info(`Syncing scrobble ${chalk.cyan(scrobble.id)} ...`); - try { - await publishScrobble(ctx, scrobble.id); - } catch (err) { - consola.error( - `Failed to sync scrobble ${chalk.cyan(scrobble.id)}:`, - err, - ); - } - } - } - process.exit(0); -} - -for (const arg of args) { - consola.info(`Syncing scrobbles ${chalk.magenta(arg)} ...`); - await updateUris(arg); - +async function refreshScrobbles(ctx: Context, did: string) { const records = await ctx.db .select({ - scrobble: scrobbles, + scrobble: tables.scrobbles, }) - .from(scrobbles) - .innerJoin(users, eq(scrobbles.userId, users.id)) - .where(or(eq(users.did, arg), eq(users.handle, arg))) - .orderBy(desc(scrobbles.createdAt)) - .limit(process.env.SYNC_SIZE ? parseInt(process.env.SYNC_SIZE, 10) : 20); + .from(tables.scrobbles) + .innerJoin(tables.users, eq(tables.scrobbles.userId, tables.users.id)) + .where(or(eq(tables.users.did, did), eq(tables.users.handle, did))) + .orderBy(desc(tables.scrobbles.createdAt)) + .limit(5); for (const { scrobble } of records) { consola.info(`Syncing scrobble ${chalk.cyan(scrobble.id)} ...`); @@ -203,7 +159,3 @@ for (const arg of args) { } consola.info(`Synced ${chalk.greenBright(records.length)} scrobbles`); } - -process.exit(0); - -*/ diff --git a/crates/jetstream/src/repo.rs b/crates/jetstream/src/repo.rs index 6a61002e..f59ae5f8 100644 --- a/crates/jetstream/src/repo.rs +++ b/crates/jetstream/src/repo.rs @@ -23,8 +23,8 @@ use crate::{ webhook_worker::{push_to_queue, AppState}, xata::{ album::Album, album_track::AlbumTrack, artist::Artist, artist_album::ArtistAlbum, - artist_track::ArtistTrack, track::Track, user::User, user_album::UserAlbum, - user_artist::UserArtist, user_track::UserTrack, + artist_track::ArtistTrack, scrobble::Scrobble, track::Track, user::User, + user_album::UserAlbum, user_artist::UserArtist, user_track::UserTrack, }, }; @@ -98,6 +98,19 @@ pub async fn save_scrobble( .await?; tx.commit().await?; + + let scrobbles: Vec = sqlx::query_as::<_, Scrobble>( + r#" + SELECT * FROM scrobbles WHERE user_id = $1 ORDER BY created_at DESC LIMIT 5 + "#, + ) + .bind(&user_id) + .fetch_all(&*pool) + .await?; + + let scrobble_id = scrobbles[0].xata_id.clone(); + nc.publish("rocksky.scrobble.new", scrobble_id.into()) + .await?; publish_user(&nc, &pool, &user_id).await?; let users: Vec = -- 2.51.2