diff --git a/src/backend/common/vendor/RockSkyApiClient.ts b/src/backend/common/vendor/RockSkyApiClient.ts index cb52c6c0..cf2aa437 100644 --- a/src/backend/common/vendor/RockSkyApiClient.ts +++ b/src/backend/common/vendor/RockSkyApiClient.ts @@ -20,6 +20,7 @@ import { CreateScrobbleInput, RockskyClient } from "@rocksky/sdk"; import { getRoot } from "../../ioc.js"; import { MSCache } from "../Cache.js"; import { HandleData } from "../infrastructure/config/client/atproto.js"; +import { parseRegexSingle } from "@foxxmd/regex-buddy-core"; interface SubmitOptions { log?: boolean @@ -73,7 +74,7 @@ export class RockSkyApiClient extends AbstractApiClient { this.rsClient = new RockskyClient({auth: token}); } - protected isLzMode = () => this.config.key !== undefined; + isLzMode = () => this.config.key !== undefined; doCallLZApi = async (req: Request, retries = 0): Promise => { try { @@ -247,7 +248,12 @@ interface UserScrobbleResponse { scrobbles?: RockskyScrobble[] } -export const rockskyScrobbleToPlay = (obj: RockskyScrobble): PlayObject => { +export const rockskyScrobbleToPlay = (obj: RockskyScrobble, opts: {playId?: string, web?: string, user?: string} = {}): PlayObject => { + const { + playId, + web, + user + } = opts; const play: PlayObjectLifecycleless = { data: { track: obj.title, @@ -259,15 +265,40 @@ export const rockskyScrobbleToPlay = (obj: RockskyScrobble): PlayObject => { playDate: dayjs.utc(obj.createdAt).local() }, meta: { - // @ts-expect-error its in the response but missing from types - trackId: obj.trackId, - playId: obj.id, + playId, + user } }; + if(web !== undefined) { + play.meta.url = {web}; + } + if('trackId' in obj) { + play.meta.trackId = obj.trackId as string; + } + // if('id' in obj) { + // play.meta.playId = obj.id as string; + // } if('albumArt' in obj) { play.meta.art = {album: obj.albumArt as string} } + if(obj.uri !== undefined) { + const uriRes = parseRegexSingle(ATPROTO_URI_REGEX, obj.uri); + if(uriRes !== undefined) { + if(web === undefined) { + play.meta.url = { + web: `https://atproto.at/viewer?uri=${uriRes.named.resource}` + } + } + if(playId === undefined) { + play.meta.playId = uriRes.named.tid; + } + if(user === undefined) { + play.meta.user = uriRes.named.did; + } + } + } + return baseFormatPlayObj(obj, play); } @@ -283,4 +314,6 @@ export const playToRockskyRecord = (play: PlayObject): CreateScrobbleInput => { timestamp: play.data.playDate.unix() } return csi; -} \ No newline at end of file +} + +const ATPROTO_URI_REGEX = new RegExp(/at:\/\/(?(?did.*?)\/app\.rocksky\.scrobble\/(?.*))/); \ No newline at end of file diff --git a/src/backend/common/vendor/atproto/ATProtoUnauthenticatedApiClient.ts b/src/backend/common/vendor/atproto/ATProtoUnauthenticatedApiClient.ts index 95386a6e..7eb2f328 100644 --- a/src/backend/common/vendor/atproto/ATProtoUnauthenticatedApiClient.ts +++ b/src/backend/common/vendor/atproto/ATProtoUnauthenticatedApiClient.ts @@ -11,7 +11,7 @@ export class ATProtoUnauthenticatedApiClient extends AbstractATProtoApiClient { declare client: Client; async initClient(): Promise { - this.userData = await getATProtoIdentifier(this.config, {logger: this.logger, cache: this.cache.cacheAuth}); + await this.hydrateHandleData(); this.client = new Client({ handler: simpleFetchHandler({ service: this.userData.pds }) }); } diff --git a/src/backend/common/vendor/atproto/AbstractATProtoApiClient.ts b/src/backend/common/vendor/atproto/AbstractATProtoApiClient.ts index 3035647d..298eebe3 100644 --- a/src/backend/common/vendor/atproto/AbstractATProtoApiClient.ts +++ b/src/backend/common/vendor/atproto/AbstractATProtoApiClient.ts @@ -5,7 +5,7 @@ import { MSCache } from "../../Cache.js"; import { UpstreamError } from "../../errors/UpstreamError.js"; import { streamBodyProgress } from "../../../utils/NetworkUtils.js"; import { ATProtoUserIdentifierData, HandleData } from "../../infrastructure/config/client/atproto.js"; -import { checkPds, isDID, identifierToAtProtoHandle } from "./atUtils.js"; +import { checkPds, isDID, identifierToAtProtoHandle, getATProtoIdentifier } from "./atUtils.js"; import { Client, isXRPCErrorPayload } from '@atcute/client'; import { ComAtprotoSyncGetRepo } from '@atcute/atproto'; import { AtprotoDid } from "@atcute/lexicons/syntax"; @@ -20,16 +20,22 @@ export abstract class AbstractATProtoApiClient extends AbstractApiClient { cache: MSCache; - constructor(name: any, config: ATProtoUserIdentifierData, options: AbstractApiOptions) { + constructor(name: any, config: ATProtoUserIdentifierData & {handleData?: HandleData}, options: AbstractApiOptions) { super('atproto', name, config, options); this.cache = getRoot().items.cache(); - const cleanIdentifier = this.config.identifier; - if(isDID(cleanIdentifier)) { - this.logger.debug(`Identifier ${cleanIdentifier} looks like a DID, skipping parsing as a handle.`); - this.config.did = cleanIdentifier; + if(config.handleData !== undefined) { + this.userData = config.handleData; + this.config.did = config.handleData.did; + this.config.identifier = config.handleData.handle; } else { - this.config.identifier = identifierToAtProtoHandle(this.config.identifier, {logger: this.logger, defaultDomain: 'bsky.social'}); + const cleanIdentifier = this.config.identifier; + if(isDID(cleanIdentifier)) { + this.logger.debug(`Identifier ${cleanIdentifier} looks like a DID, skipping parsing as a handle.`); + this.config.did = cleanIdentifier; + } else { + this.config.identifier = identifierToAtProtoHandle(this.config.identifier, {logger: this.logger, defaultDomain: 'bsky.social'}); + } } } @@ -39,6 +45,12 @@ export abstract class AbstractATProtoApiClient extends AbstractApiClient { return await checkPds(data, {logger: this.logger, cache: this.cache.cacheAuth}); } + async hydrateHandleData(): Promise { + if(this.userData === undefined) { + this.userData = await getATProtoIdentifier(this.config, {logger: this.logger, cache: this.cache.cacheAuth}); + } + } + async getCAR(did: AtprotoDid) { const resp = await this.client.call(ComAtprotoSyncGetRepo, { params: { diff --git a/src/backend/scrobblers/RockskyScrobbler.ts b/src/backend/scrobblers/RockskyScrobbler.ts index aebaf3e8..ec50e6f7 100644 --- a/src/backend/scrobblers/RockskyScrobbler.ts +++ b/src/backend/scrobblers/RockskyScrobbler.ts @@ -1,21 +1,30 @@ -import { Logger } from "@foxxmd/logging"; +import { childLogger, Logger } from "@foxxmd/logging"; import EventEmitter from "events"; import { PlayObject, SourcePlayerObj } from "../../core/Atomic.js"; import { buildTrackString, capitalize } from "../../core/StringUtils.js"; import { isNodeNetworkException } from "../common/errors/NodeErrors.js"; -import { FormatPlayObjectOptions } from "../common/infrastructure/Atomic.js"; +import { FormatPlayObjectOptions, InternalConfigOptional } from "../common/infrastructure/Atomic.js"; import { ListenbrainzApiClient } from "../common/vendor/ListenbrainzApiClient.js"; import { playToListenPayload } from '../common/vendor/listenbrainz/lzUtils.js'; import { ListenPayload } from '../common/vendor/listenbrainz/interfaces.js'; import { Notifiers } from "../notifier/Notifiers.js"; -import AbstractScrobbleClient from "./AbstractScrobbleClient.js"; -import { isDebugMode } from "../utils.js"; -import { RockSkyApiClient, SubmitResponse } from "../common/vendor/RockSkyApiClient.js"; +import { durationToHuman, isDebugMode } from "../utils.js"; +import { RockSkyApiClient, rockskyScrobbleToPlay, SubmitResponse } from "../common/vendor/RockSkyApiClient.js"; import { RockSkyClientConfig } from "../common/infrastructure/config/client/rocksky.js"; import { ScrobbleSubmitError } from "../common/errors/MSErrors.js"; - -export default class RockskyScrobbler extends AbstractScrobbleClient { +import AbstractHistoricalScrobbleClient from "./AbstractHistoricalScrobbleClient.js"; +import { fromStream } from '@atcute/repo'; +import fsPromise from 'node:fs/promises'; +import fs from 'node:fs'; +import path from 'path'; +import dayjs from "dayjs"; +import { Readable } from 'stream'; +import { ATProtoUnauthenticatedApiClient } from "../common/vendor/atproto/ATProtoUnauthenticatedApiClient.js"; +import { playToRepositoryCreatePlayHistoricalOpts, RepositoryCreatePlayHistoricalOpts } from "../common/database/drizzle/repositories/PlayHistoricalRepository.js"; +import { isAbortError } from "abort-controller-x"; + +export default class RockskyScrobbler extends AbstractHistoricalScrobbleClient { api: RockSkyApiClient; requiresAuth = true; @@ -23,15 +32,18 @@ export default class RockskyScrobbler extends AbstractScrobbleClient { declare config: RockSkyClientConfig; - constructor(name: any, config: RockSkyClientConfig, options = {}, notifier: Notifiers, emitter: EventEmitter, logger: Logger) { + protected configDir: string; + + constructor(name: any, config: RockSkyClientConfig, options: InternalConfigOptional & { [key: string]: any }, notifier: Notifiers, emitter: EventEmitter, logger: Logger) { super('rocksky', name, config, notifier, emitter, logger); - this.api = new RockSkyApiClient(name, {...config.data, ...config.options}, {logger: this.logger}); + this.api = new RockSkyApiClient(name, { ...config.data, ...config.options }, { logger: this.logger }); // https://listenbrainz.readthedocs.io/en/latest/users/api/core.html#get--1-user-(user_name)-listens // 1000 is way too high. maxing at 100 this.MAX_INITIAL_SCROBBLES_FETCH = 100; this.supportsNowPlaying = false; // PDS rate limit for operations is ~2/sec this.scrobbleDelay = 2000; + this.configDir = options.configDir; } formatPlayObj = (obj: any, options: FormatPlayObjectOptions = {}) => ListenbrainzApiClient.formatPlayObj(obj, options); @@ -59,7 +71,7 @@ export default class RockskyScrobbler extends AbstractScrobbleClient { try { return await this.api.testAuth(); } catch (e) { - if(isNodeNetworkException(e)) { + if (isNodeNetworkException(e)) { this.logger.error('Could not communicate with Rocksky API'); } throw e; @@ -83,10 +95,10 @@ export default class RockskyScrobbler extends AbstractScrobbleClient { } = playObj; try { - const result = await this.api.submitListen(playObj, { log: isDebugMode()}); + const result = await this.api.submitListen(playObj, { log: isDebugMode() }); - if(((result.response as SubmitResponse).payload?.ignored_listens ?? 0) > 0) { - throw new ScrobbleSubmitError('Scrobble was successfully submitted but Rocksky ignored it', {showStopper: false, responseBody: result.response, payload: result.payload}); + if (this.api.isLzMode() && ((result.response as SubmitResponse).payload?.ignored_listens ?? 0) > 0) { + throw new ScrobbleSubmitError('Scrobble was successfully submitted but Rocksky ignored it', { showStopper: false, responseBody: result.response, payload: result.payload }); } if (newFromSource) { @@ -96,16 +108,141 @@ export default class RockskyScrobbler extends AbstractScrobbleClient { } return result; } catch (e) { - await this.notifier.notify({title: `Client - ${capitalize(this.type)} - ${this.name} - Scrobble Error`, message: `Failed to scrobble => ${buildTrackString(playObj)} | Error: ${e.message}`, priority: 'error'}); + await this.notifier.notify({ title: `Client - ${capitalize(this.type)} - ${this.name} - Scrobble Error`, message: `Failed to scrobble => ${buildTrackString(playObj)} | Error: ${e.message}`, priority: 'error' }); throw e; } } doPlayingNow = async (data: SourcePlayerObj) => { try { - await this.api.submitListen(data.play, { listenType: 'playing_now'}); + await this.api.submitListen(data.play, { listenType: 'playing_now' }); } catch (e) { throw e; } } + + protected async doHydrateHistoricalScrobbles(opts: {allowFailures?: boolean, signal?: AbortSignal } = {}) { + const logger = childLogger(this.logger, ['Historical Plays']); + const { + allowFailures = false, + signal + } = opts; + let file: string; + try { + logger.verbose('Fetching scrobbles from PDS...'); + file = await this.fetchCarToFile(); + signal?.throwIfAborted(); + } catch (e) { + throw new Error('Failed to fetch repo CAR', {cause: e}); + } + + try { + await this.parseScrobblesFromCar(file, 100, {allowFailures, logger: logger, signal}); + } catch (e) { + throw new Error('Failed to convert CAR without any error', {cause: e}); + } finally { + await fsPromise.rm(file); + } + } + + async fetchCarToFile() { + // TODO use `since` to get CAR diff instead of entire repo + // can use last import date from migrations table + const filename = path.resolve(this.configDir, `${this.getSafeExternalId()}-${dayjs().unix()}.car`); + const atClient = new ATProtoUnauthenticatedApiClient('rocksky', { handleData: this.api.userData, identifier: this.api.userData.handle }, { logger: this.logger }); + await atClient.initClient(); + await fsPromise.writeFile(filename, Buffer.from(await atClient.getCAR(this.api.userData.did))); + return filename; + } + + async parseScrobblesFromCar(filename: string, batchSize: number, opts: { allowFailures?: boolean, logger?: Logger, signal?: AbortSignal } = {}) { + + const { + allowFailures = false, + logger = this.logger, + signal + } = opts; + + const stream = Readable.toWeb(fs.createReadStream(filename)); + + await using repo = fromStream(stream); + + let batch: RepositoryCreatePlayHistoricalOpts[] = []; + let allGood = true; + let count = 0; + let persisted = 0; + const start = dayjs(); + + logger.info('Starting CAR conversion to historical plays...'); + + for await (const entry of repo) { + if (entry.collection === 'app.rocksky.scrobble') { + let play: PlayObject; + try { + play = rockskyScrobbleToPlay(entry.record, {user: this.api.userData.did, playId: entry.rkey, web: `${this.api.userData.did}/app.rocksky.scrobble/${entry.rkey}`}) + if (isDebugMode()) { + logger.trace(`(${count}) rKey ${entry.rkey} => ${buildTrackString(play)}`); + } + count++; + if (count % (batchSize * 5) === 0) { + logger.debug(`Processed ${count} records`); + signal?.throwIfAborted(); + } + } catch (e) { + if (isAbortError(e)) { + throw e; + } + if (allowFailures) { + this.logger.warn(new Error(`Failed to convert record ${entry.rkey} to Play but will continue`, { cause: e })); + continue; + } else { + throw new Error(`Failed to convert record ${entry.rkey} to Play`, { cause: e }); + } + } + + const existing = await this.playsHistoricalRepo.hasByUid(entry.rkey); + if (!existing) { + batch.push(playToRepositoryCreatePlayHistoricalOpts({ play })); + } + if (batch.length >= batchSize) { + try { + const [res, valid] = await this.createHistoricalPlays(batch, opts); + persisted += valid; + if (!res) { + allGood = false; + } + } catch (e) { + throw e; + } + batch = []; + } + } + } + + logger.debug('Reached end of CAR file'); + if (batch.length > 0) { + logger.debug(`Persisting remaining ${batch.length} records...`); + try { + const [res, valid] = await this.createHistoricalPlays(batch, opts); + persisted += valid; + if (!res) { + allGood = false; + } + } catch (e) { + throw e; + } + } + logger.info(`Completed CAR conversion: Result ${allGood ? 'OK' : 'Some Errors'} in ${durationToHuman(dayjs.duration(dayjs().diff(start)))} | Records ${count} | Persisted ${persisted}`) + } + + protected async syncRecentHistoricalScrobbles(): Promise { + const recentPlays = await this.getScrobblesForTimeRange(undefined); + const unseenPlays: PlayObject[] = []; + for (const p of recentPlays) { + if(!(await this.playsHistoricalRepo.hasByUid(p.meta.playId))) { + unseenPlays.push(p); + } + } + return unseenPlays; + } } diff --git a/src/backend/scrobblers/ScrobbleClients.ts b/src/backend/scrobblers/ScrobbleClients.ts index 36a0daa7..aa6e741a 100644 --- a/src/backend/scrobblers/ScrobbleClients.ts +++ b/src/backend/scrobblers/ScrobbleClients.ts @@ -468,7 +468,7 @@ ${sources.join('\n')}`); break; case 'rocksky': const RockskyScrobbler = (await import('./RockskyScrobbler.js')).default; - newClient = new RockskyScrobbler(name, {...clientConfig, data: {configDir: this.internalConfig.configDir, ...d}, options: compositeOptions } as unknown as RockSkyClientConfig, {}, notifier, this.emitter, this.logger); + newClient = new RockskyScrobbler(name, {...clientConfig, data: {configDir: this.internalConfig.configDir, ...d}, options: compositeOptions } as unknown as RockSkyClientConfig, this.internalConfig, notifier, this.emitter, this.logger); break; case 'discord': const DiscordScrobbler = (await import('./DiscordScrobbler.js')).default;