A decentralized music tracking and discovery platform built on AT Protocol 🎵 forked from rocksky.app/rocksky
Something went wrong. Try again.
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084import { JetStreamClient, JetStreamEvent } from "jetstream";import { logger } from "logger";import { ctx } from "context";import { Agent } from "@atproto/api";import { env } from "lib/env";import { createAgent } from "lib/agent";import chalk from "chalk";import * as Artist from "lexicon/types/app/rocksky/artist";import * as Album from "lexicon/types/app/rocksky/album";import * as Song from "lexicon/types/app/rocksky/song";import * as Scrobble from "lexicon/types/app/rocksky/scrobble";import { SelectUser } from "schema/users";import schema from "schema";import { createId } from "@paralleldrive/cuid2";import _ from "lodash";import { and, eq, or } from "drizzle-orm";import { indexBy } from "ramda";import fs from "node:fs";import os from "node:os";import path from "node:path";import { getDidAndHandle } from "lib/getDidAndHandle";import { cleanUpJetstreamLockOnExit } from "lib/cleanUpJetstreamLock";import { cleanUpSyncLockOnExit } from "lib/cleanUpSyncLock";import { CarReader } from "@ipld/car";import * as cbor from "@ipld/dag-cbor";
type Artists = { value: Artist.Record; uri: string; cid: string }[];type Albums = { value: Album.Record; uri: string; cid: string }[];type Songs = { value: Song.Record; uri: string; cid: string }[];type Scrobbles = { value: Scrobble.Record; uri: string; cid: string }[];
export async function sync() { const [did, handle] = await getDidAndHandle(); const agent: Agent = await createAgent(did, handle);
const user = await createUser(agent, did, handle); await subscribeToJetstream(user);
logger.info` DID: ${did}`; logger.info` Handle: ${handle}`;
const carReader = await downloadCarFile(agent);
const [artists, albums, songs, scrobbles] = await Promise.all([ getRockskyUserArtists(agent, carReader), getRockskyUserAlbums(agent, carReader), getRockskyUserSongs(agent, carReader), getRockskyUserScrobbles(agent, carReader), ]);
logger.info` Artists: ${artists.length}`; logger.info` Albums: ${albums.length}`; logger.info` Songs: ${songs.length}`; logger.info` Scrobbles: ${scrobbles.length}`;
const lockFilePath = path.join(os.tmpdir(), `rocksky-${did}.lock`);
if (await fs.promises.stat(lockFilePath).catch(() => false)) { logger.error`Lock file already exists, if you want to force sync, delete the lock file ${lockFilePath}`; process.exit(1); }
await fs.promises.writeFile(lockFilePath, ""); cleanUpSyncLockOnExit(user.did);
await createArtists(artists, user); await createAlbums(albums, user); await createSongs(songs, user); await createScrobbles(scrobbles, user);
await fs.promises.unlink(lockFilePath);}
const getEndpoint = () => { const endpoint = env.JETSTREAM_SERVER;
if (endpoint?.endsWith("/subscribe")) { return endpoint; }
return `${endpoint}/subscribe`;};
export const createUser = async ( agent: Agent, did: string, handle: string,): Promise<SelectUser> => { const { data: profileRecord } = await agent.com.atproto.repo.getRecord({ repo: agent.assertDid, collection: "app.bsky.actor.profile", rkey: "self", });
const displayName = _.get(profileRecord, "value.displayName") as | string | undefined; const avatar = `https://cdn.bsky.app/img/avatar/plain/${did}/${_.get(profileRecord, "value.avatar.ref", "").toString()}@jpeg`;
const [user] = await ctx.db .insert(schema.users) .values({ id: createId(), did, handle, displayName, avatar, }) .onConflictDoUpdate({ target: schema.users.did, set: { handle, displayName, avatar, updatedAt: new Date(), }, }) .returning() .execute();
return user;};
const createArtists = async (artists: Artists, user: SelectUser) => { if (artists.length === 0) return;
const tags = artists.map((artist) => artist.value.tags || []);
// Batch genre inserts to avoid stack overflow const uniqueTags = tags .flat() .filter((tag) => tag) .map((tag) => ({ id: createId(), name: tag, }));
const BATCH_SIZE = 1000; for (let i = 0; i < uniqueTags.length; i += BATCH_SIZE) { const batch = uniqueTags.slice(i, i + BATCH_SIZE); await ctx.db .insert(schema.genres) .values(batch) .onConflictDoNothing({ target: schema.genres.name, }) .execute(); }
const genres = await ctx.db.select().from(schema.genres).execute();
const genreMap = indexBy((genre) => genre.name, genres);
// Process artists in batches let totalArtistsImported = 0;
for (let i = 0; i < artists.length; i += BATCH_SIZE) { const batch = artists.slice(i, i + BATCH_SIZE);
ctx.db.transaction((tx) => { const newArtists = tx .insert(schema.artists) .values( batch.map((artist) => ({ id: createId(), name: artist.value.name, cid: artist.cid, uri: artist.uri, biography: artist.value.bio, born: artist.value.born ? new Date(artist.value.born) : null, bornIn: artist.value.bornIn, died: artist.value.died ? new Date(artist.value.died) : null, picture: artist.value.pictureUrl, genres: artist.value.tags?.join(", "), })), ) .onConflictDoNothing({ target: schema.artists.cid, }) .returning() .all();
if (newArtists.length === 0) return;
const artistGenres = newArtists .map( (artist) => artist.genres ?.split(", ") .filter((tag) => !!tag && !!genreMap[tag]) .map((tag) => ({ id: createId(), artistId: artist.id, genreId: genreMap[tag].id, })) || [], ) .flat();
if (artistGenres.length > 0) { tx.insert(schema.artistGenres) .values(artistGenres) .onConflictDoNothing({ target: [schema.artistGenres.artistId, schema.artistGenres.genreId], }) .returning() .run(); }
tx.insert(schema.userArtists) .values( newArtists.map((artist) => ({ id: createId(), userId: user.id, artistId: artist.id, uri: artist.uri, })), ) .run();
totalArtistsImported += newArtists.length; }); }
logger.info`👤 ${totalArtistsImported} Artists imported`;};
const createAlbums = async (albums: Albums, user: SelectUser) => { if (albums.length === 0) return;
const artists = await Promise.all( albums.map(async (album) => ctx.db .select() .from(schema.artists) .where(eq(schema.artists.name, album.value.artist)) .execute() .then(([artist]) => artist), ), );
const validAlbumData = albums .map((album, index) => ({ album, artist: artists[index] })) .filter(({ artist }) => artist);
// Process albums in batches const BATCH_SIZE = 1000; let totalAlbumsImported = 0;
for (let i = 0; i < validAlbumData.length; i += BATCH_SIZE) { const batch = validAlbumData.slice(i, i + BATCH_SIZE);
ctx.db.transaction((tx) => { const newAlbums = tx .insert(schema.albums) .values( batch.map(({ album, artist }) => ({ id: createId(), cid: album.cid, uri: album.uri, title: album.value.title, artist: album.value.artist, releaseDate: album.value.releaseDate, year: album.value.year, albumArt: album.value.albumArtUrl, artistUri: artist.uri, appleMusicLink: album.value.appleMusicLink, spotifyLink: album.value.spotifyLink, tidalLink: album.value.tidalLink, youtubeLink: album.value.youtubeLink, })), ) .onConflictDoNothing({ target: schema.albums.cid, }) .returning() .all();
if (newAlbums.length === 0) return;
tx.insert(schema.userAlbums) .values( newAlbums.map((album) => ({ id: createId(), userId: user.id, albumId: album.id, uri: album.uri, })), ) .run();
totalAlbumsImported += newAlbums.length; }); }
logger.info`💿 ${totalAlbumsImported} Albums imported`;};
const createSongs = async (songs: Songs, user: SelectUser) => { if (songs.length === 0) return;
const albums = await Promise.all( songs.map((song) => ctx.db .select() .from(schema.albums) .where( and( eq(schema.albums.artist, song.value.albumArtist), eq(schema.albums.title, song.value.album), ), ) .execute() .then((result) => result[0]), ), );
const artists = await Promise.all( songs.map((song) => ctx.db .select() .from(schema.artists) .where(eq(schema.artists.name, song.value.albumArtist)) .execute() .then((result) => result[0]), ), );
const validSongData = songs .map((song, index) => ({ song, artist: artists[index], album: albums[index], })) .filter(({ artist, album }) => artist && album);
// Process in batches to avoid stack overflow with large datasets const BATCH_SIZE = 1000; let totalTracksImported = 0;
for (let i = 0; i < validSongData.length; i += BATCH_SIZE) { const batch = validSongData.slice(i, i + BATCH_SIZE); const batchNumber = Math.floor(i / BATCH_SIZE) + 1; const totalBatches = Math.ceil(validSongData.length / BATCH_SIZE);
logger.info`▶️ Processing tracks batch ${batchNumber}/${totalBatches} (${Math.min(i + BATCH_SIZE, validSongData.length)}/${validSongData.length})`;
ctx.db.transaction((tx) => { const tracks = tx .insert(schema.tracks) .values( batch.map(({ song, artist, album }) => ({ id: createId(), cid: song.cid, uri: song.uri, title: song.value.title, artist: song.value.artist, albumArtist: song.value.albumArtist, albumArt: song.value.albumArtUrl, album: song.value.album, trackNumber: song.value.trackNumber, duration: song.value.duration, mbId: song.value.mbid, youtubeLink: song.value.youtubeLink, spotifyLink: song.value.spotifyLink, appleMusicLink: song.value.appleMusicLink, tidalLink: song.value.tidalLink, discNumber: song.value.discNumber, lyrics: song.value.lyrics, composer: song.value.composer, genre: song.value.genre, label: song.value.label, copyrightMessage: song.value.copyrightMessage, albumUri: album.uri, artistUri: artist.uri, })), ) .onConflictDoNothing() .returning() .all();
if (tracks.length === 0) return;
tx.insert(schema.albumTracks) .values( tracks.map((track, index) => ({ id: createId(), albumId: batch[index].album.id, trackId: track.id, })), ) .onConflictDoNothing({ target: [schema.albumTracks.albumId, schema.albumTracks.trackId], }) .run();
tx.insert(schema.userTracks) .values( tracks.map((track) => ({ id: createId(), userId: user.id, trackId: track.id, uri: track.uri, })), ) .onConflictDoNothing({ target: [schema.userTracks.userId, schema.userTracks.trackId], }) .run();
totalTracksImported += tracks.length; }); }
logger.info`▶️ ${totalTracksImported} Tracks imported`;};
const createScrobbles = async (scrobbles: Scrobbles, user: SelectUser) => { if (!scrobbles.length) return;
logger.info`Loading Scrobble Tracks ...`;
const tracks = await Promise.all( scrobbles.map((scrobble) => ctx.db .select() .from(schema.tracks) .where( and( eq(schema.tracks.title, scrobble.value.title), eq(schema.tracks.artist, scrobble.value.artist), eq(schema.tracks.album, scrobble.value.album), eq(schema.tracks.albumArtist, scrobble.value.albumArtist), ), ) .execute() .then(([track]) => track), ), );
logger.info`Loading Scrobble Albums ...`;
const albums = await Promise.all( scrobbles.map((scrobble) => ctx.db .select() .from(schema.albums) .where( and( eq(schema.albums.title, scrobble.value.album), eq(schema.albums.artist, scrobble.value.albumArtist), ), ) .execute() .then(([album]) => album), ), );
logger.info`Loading Scrobble Artists ...`;
const artists = await Promise.all( scrobbles.map((scrobble) => ctx.db .select() .from(schema.artists) .where( or( and(eq(schema.artists.name, scrobble.value.artist)), and(eq(schema.artists.name, scrobble.value.albumArtist)), ), ) .execute() .then(([artist]) => artist), ), );
const validScrobbleData = scrobbles .map((scrobble, index) => ({ scrobble, track: tracks[index], album: albums[index], artist: artists[index], })) .filter(({ track, album, artist }) => track && album && artist);
// Process in batches to avoid stack overflow with large datasets const BATCH_SIZE = 1000; let totalScrobblesImported = 0;
for (let i = 0; i < validScrobbleData.length; i += BATCH_SIZE) { const batch = validScrobbleData.slice(i, i + BATCH_SIZE); const batchNumber = Math.floor(i / BATCH_SIZE) + 1; const totalBatches = Math.ceil(validScrobbleData.length / BATCH_SIZE);
logger.info`🕒 Processing scrobbles batch ${batchNumber}/${totalBatches} (${Math.min(i + BATCH_SIZE, validScrobbleData.length)}/${validScrobbleData.length})`;
const result = await ctx.db .insert(schema.scrobbles) .values( batch.map(({ scrobble, track, album, artist }) => ({ id: createId(), userId: user.id, trackId: track.id, albumId: album.id, artistId: artist.id, uri: scrobble.uri, cid: scrobble.cid, timestamp: new Date(scrobble.value.createdAt), })), ) .onConflictDoNothing({ target: schema.scrobbles.cid, }) .returning() .execute();
totalScrobblesImported += result.length; }
logger.info`🕒 ${totalScrobblesImported} scrobbles imported`;};
export const subscribeToJetstream = (user: SelectUser): Promise<void> => { const lockFile = path.join(os.tmpdir(), `rocksky-jetstream-${user.did}.lock`); if (fs.existsSync(lockFile)) { logger.warn`JetStream subscription already in progress for user ${user.did}`; logger.warn`Skipping subscription`; logger.warn`Lock file exists at ${lockFile}`; return Promise.resolve(); }
fs.writeFileSync(lockFile, "");
const client = new JetStreamClient({ wantedCollections: [ "app.rocksky.scrobble", "app.rocksky.artist", "app.rocksky.album", "app.rocksky.song", ], endpoint: getEndpoint(), wantedDids: [user.did],
// Reconnection settings maxReconnectAttempts: 10, reconnectDelay: 1000, maxReconnectDelay: 30000, backoffMultiplier: 1.5,
// Enable debug logging debug: true, });
return new Promise((resolve, reject) => { client.on("open", () => { logger.info`✅ Connected to JetStream!`; cleanUpJetstreamLockOnExit(user.did); resolve(); });
client.on("message", async (data) => { const event = data as JetStreamEvent;
if (event.kind === "commit" && event.commit) { const { operation, collection, record, rkey, cid } = event.commit; const uri = `at://${event.did}/${collection}/${rkey}`;
logger.info`\n📡 New event:`; logger.info` Operation: ${operation}`; logger.info` Collection: ${collection}`; logger.info` DID: ${event.did}`; logger.info` Uri: ${uri}`;
if (operation === "create" && record) { console.log(JSON.stringify(record, null, 2)); await onNewCollection(record, cid, uri, user); }
logger.info` Cursor: ${event.time_us}`; } });
client.on("error", (error) => { logger.error`❌ Error: ${error}`; cleanUpJetstreamLockOnExit(user.did); reject(error); });
client.on("reconnect", (data) => { const { attempt } = data as { attempt: number }; logger.info`🔄 Reconnecting... (attempt ${attempt})`; });
client.connect(); });};
const onNewCollection = async ( record: any, cid: string, uri: string, user: SelectUser,) => { switch (record.$type) { case "app.rocksky.song": await onNewSong(record, cid, uri, user); break; case "app.rocksky.album": await onNewAlbum(record, cid, uri, user); break; case "app.rocksky.artist": await onNewArtist(record, cid, uri, user); break; case "app.rocksky.scrobble": await onNewScrobble(record, cid, uri, user); break; default: logger.warn`Unknown collection type: ${record.$type}`; }};
const onNewSong = async ( record: Song.Record, cid: string, uri: string, user: SelectUser,) => { const { title, artist, album } = record; logger.info` New song: ${title} by ${artist} from ${album}`; await createSongs( [ { cid, uri, value: record, }, ], user, );};
const onNewAlbum = async ( record: Album.Record, cid: string, uri: string, user: SelectUser,) => { const { title, artist } = record; logger.info` New album: ${title} by ${artist}`; await createAlbums( [ { cid, uri, value: record, }, ], user, );};
const onNewArtist = async ( record: Artist.Record, cid: string, uri: string, user: SelectUser,) => { const { name } = record; logger.info` New artist: ${name}`; await createArtists( [ { cid, uri, value: record, }, ], user, );};
const onNewScrobble = async ( record: Scrobble.Record, cid: string, uri: string, user: SelectUser,) => { const { title, createdAt, artist, album, albumArtist } = record; logger.info` New scrobble: ${title} at ${createdAt}`;
// Check if the artist exists, create if not let [artistRecord] = await ctx.db .select() .from(schema.artists) .where(eq(schema.artists.name, record.albumArtist)) .execute();
if (!artistRecord) { logger.info` ⚙️ Artist not found, creating: "${albumArtist}"`;
// Create a synthetic artist record from scrobble data const artistUri = `at://${user.did}/app.rocksky.artist/${createId()}`; const artistCid = createId();
await createArtists( [ { cid: artistCid, uri: artistUri, value: { $type: "app.rocksky.artist", name: record.albumArtist, createdAt: new Date().toISOString(), tags: record.tags || [], } as Artist.Record, }, ], user, );
[artistRecord] = await ctx.db .select() .from(schema.artists) .where(eq(schema.artists.name, record.albumArtist)) .execute();
if (!artistRecord) { logger.error` ❌ Failed to create artist. Skipping scrobble.`; return; } }
// Check if the album exists, create if not let [albumRecord] = await ctx.db .select() .from(schema.albums) .where( and( eq(schema.albums.title, record.album), eq(schema.albums.artist, record.albumArtist), ), ) .execute();
if (!albumRecord) { logger.info` ⚙️ Album not found, creating: "${album}" by ${albumArtist}`;
// Create a synthetic album record from scrobble data const albumUri = `at://${user.did}/app.rocksky.album/${createId()}`; const albumCid = createId();
await createAlbums( [ { cid: albumCid, uri: albumUri, value: { $type: "app.rocksky.album", title: record.album, artist: record.albumArtist, createdAt: new Date().toISOString(), releaseDate: record.releaseDate, year: record.year, albumArt: record.albumArt, artistUri: artistRecord.uri, spotifyLink: record.spotifyLink, appleMusicLink: record.appleMusicLink, tidalLink: record.tidalLink, youtubeLink: record.youtubeLink, } as Album.Record, }, ], user, );
// Fetch the newly created album [albumRecord] = await ctx.db .select() .from(schema.albums) .where( and( eq(schema.albums.title, record.album), eq(schema.albums.artist, record.albumArtist), ), ) .execute();
if (!albumRecord) { logger.error` ❌ Failed to create album. Skipping scrobble.`; return; } }
// Check if the track exists, create if not let [track] = await ctx.db .select() .from(schema.tracks) .where( and( eq(schema.tracks.title, record.title), eq(schema.tracks.artist, record.artist), eq(schema.tracks.album, record.album), eq(schema.tracks.albumArtist, record.albumArtist), ), ) .execute();
if (!track) { logger.info` ⚙️ Track not found, creating: "${title}" by ${artist} from ${album}`;
// Create a synthetic track record from scrobble data const trackUri = `at://${user.did}/app.rocksky.song/${createId()}`; const trackCid = createId();
await createSongs( [ { cid: trackCid, uri: trackUri, value: { $type: "app.rocksky.song", title: record.title, artist: record.artist, albumArtist: record.albumArtist, album: record.album, duration: record.duration, trackNumber: record.trackNumber, discNumber: record.discNumber, releaseDate: record.releaseDate, year: record.year, genre: record.genre, tags: record.tags, composer: record.composer, lyrics: record.lyrics, copyrightMessage: record.copyrightMessage, albumArt: record.albumArt, youtubeLink: record.youtubeLink, spotifyLink: record.spotifyLink, tidalLink: record.tidalLink, appleMusicLink: record.appleMusicLink, createdAt: new Date().toISOString(), mbId: record.mbid, label: record.label, albumUri: albumRecord.uri, artistUri: artistRecord.uri, } as Song.Record, }, ], user, );
// Fetch the newly created track [track] = await ctx.db .select() .from(schema.tracks) .where( and( eq(schema.tracks.title, record.title), eq(schema.tracks.artist, record.artist), eq(schema.tracks.album, record.album), eq(schema.tracks.albumArtist, record.albumArtist), ), ) .execute();
if (!track) { logger.error` ❌ Failed to create track. Skipping scrobble.`; return; } }
logger.info` ✓ All required entities ready. Creating scrobble...`;
await createScrobbles( [ { cid, uri, value: record, }, ], user, );};
const downloadCarFile = async (agent: Agent) => { logger.info(`Fetching repository CAR file ...`);
const repoRes = await agent.com.atproto.sync.getRepo({ did: agent.assertDid, });
return CarReader.fromBytes(new Uint8Array(repoRes.data));};
const getRockskyUserSongs = async ( agent: Agent, carReader: CarReader,): Promise<Songs> => { const results: { value: Song.Record; uri: string; cid: string; }[] = [];
try { const collection = "app.rocksky.song"; logger.info`Extracting ${collection} records from CAR file ...`;
for await (const { cid, bytes } of carReader.blocks()) { try { const decoded = cbor.decode(bytes);
// Check if this is a record with $type matching our collection if (decoded && typeof decoded === "object" && "$type" in decoded) { if (decoded.$type === collection) { const value = decoded as unknown as Song.Record; // Extract rkey from uri if present in the block, otherwise use cid const uri = `at://${agent.assertDid}/${collection}/${cid.toString()}`;
results.push({ value, uri, cid: cid.toString(), }); } } } catch (e) { logger.warn` Skipping block with CID ${cid.toString()} due to decode error: ${e}`; continue; } }
logger.info( `${chalk.cyanBright(agent.assertDid)} ${chalk.greenBright(results.length)} songs`, ); } catch (error) { logger.error(`Error fetching songs from CAR: ${error}`); throw error; }
return results;};
const getRockskyUserAlbums = async ( agent: Agent, carReader: CarReader,): Promise<Albums> => { const results: { value: Album.Record; uri: string; cid: string; }[] = [];
try { const collection = "app.rocksky.album"; logger.info`Extracting ${collection} records from CAR file ...`;
for await (const { cid, bytes } of carReader.blocks()) { try { const decoded = cbor.decode(bytes);
if (decoded && typeof decoded === "object" && "$type" in decoded) { if (decoded.$type === collection) { const value = decoded as unknown as Album.Record; const uri = `at://${agent.assertDid}/${collection}/${cid.toString()}`;
results.push({ value, uri, cid: cid.toString(), }); } } } catch (e) { logger.warn` Skipping block with CID ${cid.toString()} due to decode error: ${e}`; continue; } }
logger.info( `${chalk.cyanBright(agent.assertDid)} ${chalk.greenBright(results.length)} albums`, ); } catch (error) { logger.error(`Error fetching albums from CAR: ${error}`); throw error; }
return results;};
const getRockskyUserArtists = async ( agent: Agent, carReader: CarReader,): Promise<Artists> => { const results: { value: Artist.Record; uri: string; cid: string; }[] = [];
try { const collection = "app.rocksky.artist"; logger.info`Extracting ${collection} records from CAR file ...`;
for await (const { cid, bytes } of carReader.blocks()) { try { const decoded = cbor.decode(bytes);
if (decoded && typeof decoded === "object" && "$type" in decoded) { if (decoded.$type === collection) { const value = decoded as unknown as Artist.Record; const uri = `at://${agent.assertDid}/${collection}/${cid.toString()}`;
results.push({ value, uri, cid: cid.toString(), }); } } } catch (e) { // Skip blocks that can't be decoded continue; } }
logger.info( `${chalk.cyanBright(agent.assertDid)} ${chalk.greenBright(results.length)} artists`, ); } catch (error) { logger.error(`Error fetching artists from CAR: ${error}`); throw error; }
return results;};
const getRockskyUserScrobbles = async ( agent: Agent, carReader: CarReader,): Promise<Scrobbles> => { const results: { value: Scrobble.Record; uri: string; cid: string; }[] = [];
try { const collection = "app.rocksky.scrobble"; logger.info`Extracting ${collection} records from CAR file ...`;
for await (const { cid, bytes } of carReader.blocks()) { try { const decoded = cbor.decode(bytes);
if (decoded && typeof decoded === "object" && "$type" in decoded) { if (decoded.$type === collection) { const value = decoded as unknown as Scrobble.Record; const uri = `at://${agent.assertDid}/${collection}/${cid.toString()}`;
results.push({ value, uri, cid: cid.toString(), }); } } } catch (e) { logger.warn` Skipping block with CID ${cid.toString()} due to decode error: ${e}`; continue; } }
logger.info( `${chalk.cyanBright(agent.assertDid)} ${chalk.greenBright(results.length)} scrobbles`, ); } catch (error) { logger.error(`Error fetching scrobbles from CAR: ${error}`); throw error; }
return results;};