From 67f6880d32a2f9daa1019fb326cddf6c50842f47 Mon Sep 17 00:00:00 2001 From: FoxxMD Date: Thu, 29 Jan 2026 16:34:25 +0000 Subject: [PATCH] feat(scrobblers): Refactor scrobbling functions to capture transactional data Capture payload, response/body, errors, and warnings so we can store them in play lifecycle for future debugging --- src/backend/common/errors/MSErrors.ts | 10 +++++ src/backend/common/errors/UpstreamError.ts | 18 ++++++--- src/backend/common/infrastructure/Atomic.ts | 4 +- src/backend/common/vendor/LastfmApiClient.ts | 13 ++++--- .../common/vendor/ListenbrainzApiClient.ts | 25 +++++++----- src/backend/common/vendor/RockSkyApiClient.ts | 31 +++++++-------- .../bluesky/AbstractBlueSkyApiClient.ts | 24 +++++++----- .../common/vendor/koito/KoitoApiClient.ts | 11 +++--- .../common/vendor/maloja/MalojaApiClient.ts | 28 +++++++++----- .../scrobblers/AbstractScrobbleClient.ts | 38 ++++++++++++++----- src/backend/scrobblers/KoitoScrobbler.ts | 8 ++-- src/backend/scrobblers/LastfmScrobbler.ts | 10 ++--- .../scrobblers/ListenbrainzScrobbler.ts | 12 +++--- src/backend/scrobblers/MalojaScrobbler.ts | 22 +++-------- src/backend/scrobblers/RockskyScrobbler.ts | 17 +++++---- src/backend/scrobblers/TealfmScrobbler.ts | 8 ++-- src/backend/sources/SpotifySource.ts | 2 +- src/backend/tests/scrobbler/TestScrobbler.ts | 4 +- src/core/Atomic.ts | 18 ++++++++- 19 files changed, 185 insertions(+), 118 deletions(-) diff --git a/src/backend/common/errors/MSErrors.ts b/src/backend/common/errors/MSErrors.ts index af22b788..5eb8d065 100644 --- a/src/backend/common/errors/MSErrors.ts +++ b/src/backend/common/errors/MSErrors.ts @@ -1,6 +1,7 @@ import { parseRegexSingle } from "@foxxmd/regex-buddy-core"; import mergeErrorCause from 'merge-error-cause'; import { findCauseByFunc, findCauseByReference } from "../../utils/ErrorUtils.js"; +import { UpstreamError, UpstreamErrorOptions } from "./UpstreamError.js"; export abstract class NamedError extends Error { public abstract name: string; @@ -94,4 +95,13 @@ export const mergeSimpleError = (err: Error): Error => { return mergeErrorCause(structuredClone(err)); } return err; +} + +export class ScrobbleSubmitError extends UpstreamError { + name = 'Scrobble Submit Error'; + payload?: T; + constructor(message: string, options?: UpstreamErrorOptions & {payload?: T}) { + super(message, options); + this.payload = options?.payload; + } } \ No newline at end of file diff --git a/src/backend/common/errors/UpstreamError.ts b/src/backend/common/errors/UpstreamError.ts index 7df9d061..94e76b20 100644 --- a/src/backend/common/errors/UpstreamError.ts +++ b/src/backend/common/errors/UpstreamError.ts @@ -2,20 +2,28 @@ import { Response } from 'superagent'; import { findCauseByFunc } from "../../utils/ErrorUtils.js"; -export class UpstreamError extends Error { +export type UpstreamErrorOptions = ErrorOptions & { showStopper?: boolean, response?: Response, responseBody?: object | string }; + +export class UpstreamError extends Error { showStopper: boolean = false; response?: Response + responseBody?: object | string - constructor(message: string, options?: { cause?: T, showStopper?: boolean, response?: Response } | undefined) { + constructor(message: string, options?: UpstreamErrorOptions | undefined) { super(message, options); - const {showStopper = false, response} = options; + const {showStopper = false, response, responseBody} = options; this.showStopper = showStopper; this.response = response; + this.responseBody = responseBody; } } export const hasUpstreamError = (err: any, showStopping?: boolean): boolean => { + return findUpstreamError(err, showStopping) !== undefined; +} + +export const findUpstreamError = (err: any, showStopping?: boolean): UpstreamError | undefined => { return findCauseByFunc(err, (e) => { if (e instanceof UpstreamError) { if (showStopping === undefined) { @@ -25,5 +33,5 @@ export const hasUpstreamError = (err: any, showStopping?: boolean): boolean => { } } return false; - }) !== undefined; -} + }); +} \ No newline at end of file diff --git a/src/backend/common/infrastructure/Atomic.ts b/src/backend/common/infrastructure/Atomic.ts index db8404de..820b537e 100644 --- a/src/backend/common/infrastructure/Atomic.ts +++ b/src/backend/common/infrastructure/Atomic.ts @@ -3,7 +3,7 @@ import { Dayjs } from "dayjs"; import { Request, Response } from "express"; import { NextFunction, ParamsDictionary, Query } from "express-serve-static-core"; import { FixedSizeList } from 'fixed-size-list'; -import { isPlayObject, PlayMeta, PlayObject } from "../../../core/Atomic.js"; +import { isPlayObject, PlayMeta, PlayObject, PlayObjectLifecycleless } from "../../../core/Atomic.js"; import TupleMap from "../TupleMap.js"; import { MusicBrainzApi } from 'musicbrainz-api'; @@ -242,7 +242,7 @@ export const SINGLE_USER_PLATFORM_ID: PlayPlatformId = [NO_DEVICE, NO_USER]; export interface ScrobbledPlayObject { play: PlayObject - scrobble: PlayObject + scrobble: PlayObjectLifecycleless } diff --git a/src/backend/common/vendor/LastfmApiClient.ts b/src/backend/common/vendor/LastfmApiClient.ts index ef097598..0d9ed1bc 100644 --- a/src/backend/common/vendor/LastfmApiClient.ts +++ b/src/backend/common/vendor/LastfmApiClient.ts @@ -1,5 +1,5 @@ import dayjs, { Dayjs } from "dayjs"; -import { BrainzMeta, PlayObject, PlayObjectLifecycleless, URLData } from "../../../core/Atomic.js"; +import { BrainzMeta, PlayObject, PlayObjectLifecycleless, ScrobbleActionResult, URLData } from "../../../core/Atomic.js"; import { nonEmptyStringOrDefault, splitByFirstFound } from "../../../core/StringUtils.js"; import { removeUndefinedKeys, sleep, writeFile } from "../../utils.js"; import { objectIsEmpty, readJson } from '../../utils/DataUtils.js'; @@ -15,6 +15,7 @@ import { LastFMUser, LastFMAuth, LastFMTrack, LastFMUserGetRecentTracksResponse, import clone from 'clone'; import { IncomingMessage } from "http"; import { baseFormatPlayObj } from "../../utils/PlayTransformUtils.js"; +import { ScrobbleSubmitError } from "../errors/MSErrors.js"; const badErrors = [ 'api key suspended', @@ -304,7 +305,7 @@ export default class LastfmApiClient extends AbstractApiClient { } } - scrobble = async (playObj: PlayObject): Promise => { + scrobble = async (playObj: PlayObject): Promise => { const { meta: { source, @@ -356,15 +357,17 @@ export default class LastfmApiClient extends AbstractApiClient { modifiedPlay.data.album = albumName; } - return modifiedPlay; + return {payload: scrobblePayload, response, mergedScrobble: modifiedPlay}; // last fm has rate limits but i can't find a specific example of what that limit is. going to default to 1 scrobble/sec to be safe //await sleep(1000); } catch (e) { + let apiError: Error; if(!(e instanceof UpstreamError)) { - throw new UpstreamError(`Error received from ${this.upstreamName} API`, {cause: e, showStopper: true}); + apiError = new UpstreamError(`Error received from ${this.upstreamName} API`, {cause: e, showStopper: true}); } else { - throw e; + apiError = e; } + throw new ScrobbleSubmitError('Failed to submit scrobble to Last.fm', {cause: apiError, payload: scrobblePayload}); } finally { this.logger.debug({payload: scrobblePayload}, 'Raw Payload'); } diff --git a/src/backend/common/vendor/ListenbrainzApiClient.ts b/src/backend/common/vendor/ListenbrainzApiClient.ts index 10ab201c..f2b5a855 100644 --- a/src/backend/common/vendor/ListenbrainzApiClient.ts +++ b/src/backend/common/vendor/ListenbrainzApiClient.ts @@ -1,7 +1,7 @@ import { stringSameness } from '@foxxmd/string-sameness'; import dayjs from "dayjs"; import request, { Request, Response } from 'superagent'; -import { BrainzMeta, PlayObject, PlayObjectLifecycleless, URLData } from "../../../core/Atomic.js"; +import { BrainzMeta, PlayObject, PlayObjectLifecycleless, ScrobbleActionResult, URLData } from "../../../core/Atomic.js"; import { combinePartsToString, slice } from "../../../core/StringUtils.js"; import { findDelimiters, @@ -24,6 +24,7 @@ import { listenObjectResponseToPlay } from './koito/KoitoApiClient.js'; import { version } from '../../ioc.js'; import { ListenPayload, ListenResponse, ListenType, MinimumTrack, SubmitListenAdditionalTrackInfo, SubmitPayload } from './listenbrainz/interfaces.js'; import { baseFormatPlayObj } from '../../utils/PlayTransformUtils.js'; +import { ScrobbleSubmitError } from '../errors/MSErrors.js'; interface SubmitOptions { log?: boolean @@ -225,13 +226,10 @@ export class ListenbrainzApiClient extends AbstractApiClient { } - submitListen = async (play: PlayObject, options: SubmitOptions = {}) => { - const { log = false, listenType = 'single'} = options; + submitListen = async (play: PlayObject, options: SubmitOptions = {}): Promise => { + const listenPayload = playToSubmitPayload(play, {listenType: options.listenType}); + const { log = false} = options; try { - const listenPayload: SubmitPayload = {listen_type: listenType, payload: [playToListenPayload(play)]}; - if(listenType === 'playing_now') { - delete listenPayload.payload[0].listened_at; - } if(log) { this.logger.debug(`Submit Payload: ${JSON.stringify(listenPayload)}`); } @@ -243,9 +241,9 @@ export class ListenbrainzApiClient extends AbstractApiClient { if(log) { this.logger.debug(`Submit Response: ${resp.text}`) } - return listenPayload; + return {payload: listenPayload, response: resp.text}; } catch (e) { - throw e; + throw new ScrobbleSubmitError(`Failed to submit to Listenbrainz (listen_type ${listenPayload.listen_type})`, {cause: e, payload: listenPayload, response: e.response, responseBody: e.response?.text}); } } @@ -819,4 +817,13 @@ export const musicServiceToCononical = (str?: string): string | undefined => { } } return undefined; +} + +export const playToSubmitPayload = (play: PlayObject, options: SubmitOptions = {}): SubmitPayload => { + const { listenType = 'single'} = options; + const listenPayload: SubmitPayload = {listen_type: listenType, payload: [playToListenPayload(play)]}; + if(listenType === 'playing_now') { + delete listenPayload.payload[0].listened_at; + } + return listenPayload; } \ No newline at end of file diff --git a/src/backend/common/vendor/RockSkyApiClient.ts b/src/backend/common/vendor/RockSkyApiClient.ts index e8072797..817c47c8 100644 --- a/src/backend/common/vendor/RockSkyApiClient.ts +++ b/src/backend/common/vendor/RockSkyApiClient.ts @@ -1,6 +1,6 @@ import dayjs from "dayjs"; import request, { Request, Response } from 'superagent'; -import { PlayObject, PlayObjectLifecycleless, URLData } from "../../../core/Atomic.js"; +import { PlayObject, PlayObjectLifecycleless, ScrobbleActionResult, URLData } from "../../../core/Atomic.js"; import { nonEmptyStringOrDefault } from "../../../core/StringUtils.js"; import { UpstreamError } from "../errors/UpstreamError.js"; import { AbstractApiOptions, DEFAULT_RETRY_MULTIPLIER, FormatPlayObjectOptions } from "../infrastructure/Atomic.js"; @@ -14,6 +14,7 @@ import { RockskyScrobble } from './rocksky/interfaces.js'; import { Handle } from "@atcute/lexicons"; import { identifierToAtProtoHandle } from './bluesky/bsUtils.js'; import { baseFormatPlayObj } from "../../utils/PlayTransformUtils.js"; +import { ScrobbleSubmitError } from "../errors/MSErrors.js"; interface SubmitOptions { log?: boolean @@ -171,21 +172,21 @@ export class RockSkyApiClient extends AbstractApiClient { } - submitListen = async (play: PlayObject, options: SubmitOptions = {}) => { + submitListen = async (play: PlayObject, options: SubmitOptions = {}): Promise => { const { log = false, listenType = 'single'} = options; - try { - - const listenPayload = playToListenPayload(play); - if(listenType === 'playing_now') { + const listenPayload = playToListenPayload(play); + if(listenType === 'playing_now') { delete listenPayload.listened_at; } - // https://tangled.org/rocksky.app/rocksky/blob/main/crates/scrobbler/src/listenbrainz/types.rs#L11 - // rocksky only uses duration_ms - if(play.data.duration !== undefined && listenPayload.track_metadata.additional_info?.duration !== undefined) { - delete listenPayload.track_metadata.additional_info.duration; - listenPayload.track_metadata.additional_info.duration_ms = Math.round(play.data.duration) * 1000; - } - const submitPayload: SubmitPayload = {listen_type: listenType, payload: [listenPayload]}; + // https://tangled.org/rocksky.app/rocksky/blob/main/crates/scrobbler/src/listenbrainz/types.rs#L11 + // rocksky only uses duration_ms + if(play.data.duration !== undefined && listenPayload.track_metadata.additional_info?.duration !== undefined) { + delete listenPayload.track_metadata.additional_info.duration; + listenPayload.track_metadata.additional_info.duration_ms = Math.round(play.data.duration) * 1000; + } + const submitPayload: SubmitPayload = {listen_type: listenType, payload: [listenPayload]}; + + try { if(log) { this.logger.debug(`Submit Payload: ${JSON.stringify(submitPayload)}`); } @@ -193,9 +194,9 @@ export class RockSkyApiClient extends AbstractApiClient { if(log) { this.logger.debug(`Submit Response: ${resp.text}`) } - return resp.body as SubmitResponse; + return {payload: submitPayload, response: resp.body as SubmitResponse}; } catch (e) { - throw e; + throw new ScrobbleSubmitError(`Error occurred while making Rocksky API scrobble (${listenType}) request`, {cause: e, payload: submitPayload}); } } diff --git a/src/backend/common/vendor/bluesky/AbstractBlueSkyApiClient.ts b/src/backend/common/vendor/bluesky/AbstractBlueSkyApiClient.ts index aa4dcae9..591902b0 100644 --- a/src/backend/common/vendor/bluesky/AbstractBlueSkyApiClient.ts +++ b/src/backend/common/vendor/bluesky/AbstractBlueSkyApiClient.ts @@ -2,9 +2,9 @@ import { getRoot } from "../../../ioc.js"; import { AbstractApiOptions } from "../../infrastructure/Atomic.js"; import { ListRecord, ScrobbleRecord, TealClientData } from "../../infrastructure/config/client/tealfm.js"; import AbstractApiClient from "../AbstractApiClient.js"; -import { Agent } from "@atproto/api"; +import { Agent, ComAtprotoRepoCreateRecord } from "@atproto/api"; import { MSCache } from "../../Cache.js"; -import { BrainzMeta, PlayObject, PlayObjectLifecycleless } from "../../../../core/Atomic.js"; +import { BrainzMeta, PlayObject, PlayObjectLifecycleless, ScrobbleActionResult } from "../../../../core/Atomic.js"; import { musicServiceToCononical } from "../ListenbrainzApiClient.js"; import { parseRegexSingle } from "@foxxmd/regex-buddy-core"; import { RecordOptions } from "../../infrastructure/config/client/tealfm.js"; @@ -12,6 +12,8 @@ import dayjs from "dayjs"; import { getScrobbleTsSOCDateWithContext } from "../../../utils/TimeUtils.js"; import { removeUndefinedKeys } from "../../../utils.js"; import { baseFormatPlayObj } from "../../../utils/PlayTransformUtils.js"; +import { ScrobbleSubmitError } from "../../errors/MSErrors.js"; +import { UpstreamError } from "../../errors/UpstreamError.js"; export abstract class AbstractBlueSkyApiClient extends AbstractApiClient { @@ -32,15 +34,17 @@ export abstract class AbstractBlueSkyApiClient extends AbstractApiClient { abstract restoreSession(): Promise; - async createScrobbleRecord(record: ScrobbleRecord): Promise { + async createScrobbleRecord(record: ScrobbleRecord): Promise { + const input: ComAtprotoRepoCreateRecord.InputSchema = { + repo: this.agent.sessionManager.did, + collection: "fm.teal.alpha.feed.play", + record + }; try { - await this.agent.com.atproto.repo.createRecord({ - repo: this.agent.sessionManager.did, - collection: "fm.teal.alpha.feed.play", - record - }); + const resp = await this.agent.com.atproto.repo.createRecord(input); + return {payload: input, response: resp.data}; } catch (e) { - throw new Error(`Failed to create record`, { cause: e }); + throw new ScrobbleSubmitError(`Failed to create record for scrobble`, { cause: e, payload: input, response: 'response' in e ? e.response : undefined }); } } @@ -53,7 +57,7 @@ export abstract class AbstractBlueSkyApiClient extends AbstractApiClient { }); return response.data.records as unknown as ListRecord[]; } catch (e) { - throw new Error(`Failed to create record`, { cause: e }); + throw new UpstreamError(`Failed to list scrobble record`, { cause: e, response: 'response' in e ? e.response : undefined }); } } } diff --git a/src/backend/common/vendor/koito/KoitoApiClient.ts b/src/backend/common/vendor/koito/KoitoApiClient.ts index 903fc1eb..a9002338 100644 --- a/src/backend/common/vendor/koito/KoitoApiClient.ts +++ b/src/backend/common/vendor/koito/KoitoApiClient.ts @@ -1,5 +1,5 @@ import dayjs from "dayjs"; -import { PlayObject, PlayObjectLifecycleless, URLData } from "../../../../core/Atomic.js"; +import { PlayObject, PlayObjectLifecycleless, ScrobbleActionResult, URLData } from "../../../../core/Atomic.js"; import { AbstractApiOptions, DEFAULT_RETRY_MULTIPLIER } from "../../infrastructure/Atomic.js"; import { KoitoData, ListenObjectResponse, ListensResponse } from "../../infrastructure/config/client/koito.js"; import AbstractApiClient from "../AbstractApiClient.js"; @@ -11,6 +11,7 @@ import { SubmitPayload } from '../listenbrainz/interfaces.js'; import { ListenType } from '../listenbrainz/interfaces.js'; import { parseRegexSingleOrFail } from "../../../utils.js"; import { baseFormatPlayObj } from "../../../utils/PlayTransformUtils.js"; +import { ScrobbleSubmitError } from "../../errors/MSErrors.js"; interface SubmitOptions { log?: boolean @@ -155,10 +156,10 @@ export class KoitoApiClient extends AbstractApiClient { } } - submitListen = async (play: PlayObject, options: SubmitOptions = {}) => { + submitListen = async (play: PlayObject, options: SubmitOptions = {}): Promise => { const { log = false, listenType = 'single' } = options; + const listenPayload: SubmitPayload = { listen_type: listenType, payload: [playToListenPayload(play)] }; try { - const listenPayload: SubmitPayload = { listen_type: listenType, payload: [playToListenPayload(play)] }; if (listenType === 'playing_now') { delete listenPayload.payload[0].listened_at; } @@ -173,9 +174,9 @@ export class KoitoApiClient extends AbstractApiClient { if (log) { this.logger.debug(`Submit Response: ${resp.text}`) } - return listenPayload; + return {payload: listenPayload, response: resp.text}; } catch (e) { - throw e; + throw new ScrobbleSubmitError(`Error occurred while making Koito API submit request (listen_type ${listenPayload.listen_type})`, {cause: e, payload: listenPayload, response: e.response, responseBody: e.response?.text}); } } diff --git a/src/backend/common/vendor/maloja/MalojaApiClient.ts b/src/backend/common/vendor/maloja/MalojaApiClient.ts index e9ca7f47..1df7028c 100644 --- a/src/backend/common/vendor/maloja/MalojaApiClient.ts +++ b/src/backend/common/vendor/maloja/MalojaApiClient.ts @@ -4,7 +4,7 @@ import compareVersions from "compare-versions"; import AbstractApiClient from "../AbstractApiClient.js"; import { getBaseFromUrl, isPortReachableConnect, joinedUrl, normalizeWebAddress } from "../../../utils/NetworkUtils.js"; import { MalojaData } from "../../infrastructure/config/client/maloja.js"; -import { PlayObject, PlayObjectLifecycleless, URLData } from "../../../../core/Atomic.js"; +import { PlayObject, PlayObjectLifecycleless, ScrobbleActionResult, URLData } from "../../../../core/Atomic.js"; import { AbstractApiOptions, DEFAULT_RETRY_MULTIPLIER, FormatPlayObjectOptions } from "../../infrastructure/Atomic.js"; import { isNodeNetworkException } from "../../errors/NodeErrors.js"; import { isSuperAgentResponseError } from "../../errors/ErrorUtils.js"; @@ -14,6 +14,7 @@ import { getMalojaResponseError, isMalojaAPIErrorBody, MalojaResponseV3CommonDat import { getScrobbleTsSOCDate, getScrobbleTsSOCDateWithContext } from '../../../utils/TimeUtils.js'; import { buildTrackString } from '../../../../core/StringUtils.js'; import { baseFormatPlayObj } from '../../../utils/PlayTransformUtils.js'; +import { ScrobbleSubmitError } from '../../errors/MSErrors.js'; @@ -188,7 +189,7 @@ export class MalojaApiClient extends AbstractApiClient { return list.map(formatPlayObj); } - scrobble = async (playObj: PlayObject): Promise<[(MalojaScrobbleData | undefined), MalojaScrobbleV3ResponseData, string?]> => { + scrobble = async (playObj: PlayObject): Promise => { const { data: { @@ -214,11 +215,10 @@ export class MalojaApiClient extends AbstractApiClient { .type('json') .send(scrobbleData)); - let scrobbleResponse: MalojaScrobbleData, - scrobbledPlay: PlayObject; - + let scrobbleResponse: MalojaScrobbleData; let responseBody: MalojaScrobbleV3ResponseData; let warnStr: string; + const msWarnings: string[] = []; responseBody = response.body; const { @@ -245,22 +245,32 @@ export class MalojaApiClient extends AbstractApiClient { ...malojaAlbum, } } + } else { + msWarnings.push('Maloja did not return track data in scrobble response! Maybe it didn\'t scrobble correctly?'); } if (warnings.length > 0) { for (const w of warnings) { warnStr = builMalojadWarningString(w); if (warnStr.includes('The submitted scrobble was not added')) { - throw new UpstreamError(`Maloja returned a warning but MS treating as error: ${warnStr}`, { showStopper: false }); + throw new ScrobbleSubmitError(`Maloja returned a warning but MS treating as error: ${warnStr}`, { showStopper: false, payload: scrobbleData, response, responseBody }); } - this.logger.warn(`Maloja Warning: ${warnStr}`); + const wstring = `Maloja Warning: ${warnStr}`; + msWarnings.push(wstring); + this.logger.warn(wstring); } } } else { - throw new UpstreamError(buildMalojaErrorString(response.body), { showStopper: false }); + throw new ScrobbleSubmitError(buildMalojaErrorString(response.body), { showStopper: false, payload: scrobbleData, response, responseBody }); } - return [scrobbleResponse, responseBody, warnStr] + return {payload: scrobbleData, warnings: msWarnings.length > 0 ? msWarnings : undefined, response: responseBody, mergedScrobble: scrobbleResponse !== undefined ? formatPlayObj(scrobbleResponse, {url: this.url.normal}) : undefined}; } catch (e) { + let scrobbleError: ScrobbleSubmitError; + if(e instanceof ScrobbleSubmitError) { + scrobbleError = e; + } else { + scrobbleError = new ScrobbleSubmitError('Error occurred while submitting scrobble to Maloja', {cause: e, payload: scrobbleData, response: 'response' in e ? e.response : undefined}); + } this.logger.error({ playInfo: buildTrackString(playObj), payload: scrobbleData }, `Scrobble Error (${sType})`); const responseError = getMalojaResponseError(e); if (responseError !== undefined) { diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index 35453d01..9bd5ae0f 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -7,13 +7,14 @@ import { MarkOptional } from "ts-essentials"; import { DeadLetterScrobble, PlayObject, - QueuedScrobble, TA_DURING, + PlayObjectLifecycleless, + QueuedScrobble, ScrobbleActionResult, ScrobblePayload, ScrobbleResponse, TA_DURING, TA_FUZZY, TrackStringOptions } from "../../core/Atomic.js"; import { buildTrackString, capitalize, truncateStringToLength } from "../../core/StringUtils.js"; import AbstractComponent from "../common/AbstractComponent.js"; -import { UpstreamError } from "../common/errors/UpstreamError.js"; +import { hasUpstreamError, UpstreamError } from "../common/errors/UpstreamError.js"; import { ARTIST_WEIGHT, Authenticatable, @@ -43,7 +44,7 @@ import { sleep, sortByOldestPlayDate, } from "../utils.js"; -import { messageWithCauses, messageWithCausesTruncatedDefault } from "../utils/ErrorUtils.js"; +import { findCauseByReference, messageWithCauses, messageWithCausesTruncatedDefault } from "../utils/ErrorUtils.js"; import { comparePlayTemporally, hasAcceptableTemporalAccuracy, @@ -60,6 +61,7 @@ import pMap, { pMapIterable } from "p-map"; import { comparePlayArtistsNormalized, comparePlayTracksNormalized } from "../utils/PlayComparisonUtils.js"; import { normalizeStr } from "../utils/StringUtils.js"; import prom, { Counter, Gauge } from 'prom-client'; +import { ScrobbleSubmitError } from "../common/errors/MSErrors.js"; type PlatformMappedPlays = Map; type NowPlayingQueue = Map; @@ -477,7 +479,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i return [validTime, log] } - addScrobbledTrack = (playObj: PlayObject, scrobbledPlay: PlayObject) => { + addScrobbledTrack = (playObj: PlayObject, scrobbledPlay: PlayObjectLifecycleless) => { this.scrobbledPlayObjs.add({play: playObj, scrobble: scrobbledPlay}); this.scrobbledCounter.labels(this.getPrometheusLabels()).inc(); this.lastScrobbledPlayDate = playObj.data.playDate; @@ -669,13 +671,20 @@ ${closestMatch.breakdowns.join('\n')}`, {leaf: ['Dupe Check']}); } } try { - return this.doScrobble(playObj); + const result = await this.doScrobble(playObj); + playObj.meta.lifecycle.scrobble = { + payload: result.payload, + warnings: result.warnings, + response: result.response, + mergedScrobble: result.mergedScrobble + } + return playObj; } finally { this.lastScrobbleAttempt = dayjs(); } } - protected abstract doScrobble(playObj: PlayObject): Promise + protected abstract doScrobble(playObj: PlayObject): Promise public abstract playToClientPayload(playObject: PlayObject): object @@ -800,10 +809,21 @@ ${closestMatch.breakdowns.join('\n')}`, {leaf: ['Dupe Check']}); try { const scrobbledPlay = await this.scrobble(transformedScrobble); this.emitEvent('scrobble', {play: transformedScrobble}); - this.addScrobbledTrack(transformedScrobble, scrobbledPlay); + this.addScrobbledTrack(scrobbledPlay, scrobbledPlay.meta.lifecycle.scrobble.mergedScrobble ?? scrobbledPlay); } catch (e) { - currQueuedPlay.play.meta.lifecycle.scrobble = this.playToClientPayload(transformedScrobble); - if (e instanceof UpstreamError && e.showStopper === false) { + currQueuedPlay.play.meta.lifecycle.scrobble = { + error: e, + payload: this.playToClientPayload(transformedScrobble) + }; + + const submitError = findCauseByReference(e, ScrobbleSubmitError); + if(submitError !== undefined) { + currQueuedPlay.play.meta.lifecycle.scrobble.error = submitError; + currQueuedPlay.play.meta.lifecycle.scrobble.payload = submitError.payload; + currQueuedPlay.play.meta.lifecycle.scrobble.response = submitError.responseBody; + } + + if (hasUpstreamError(e, false)) { this.addDeadLetterScrobble(currQueuedPlay, e); this.logger.warn(new Error(`Could not scrobble ${buildTrackString(transformedScrobble)} from Source '${currQueuedPlay.source}' but error was not show stopping. Adding scrobble to Dead Letter Queue and will retry on next heartbeat.`, {cause: e})); } else { diff --git a/src/backend/scrobblers/KoitoScrobbler.ts b/src/backend/scrobblers/KoitoScrobbler.ts index e48807d7..b96e3d2e 100644 --- a/src/backend/scrobblers/KoitoScrobbler.ts +++ b/src/backend/scrobblers/KoitoScrobbler.ts @@ -81,17 +81,17 @@ export default class KoitoScrobbler extends AbstractScrobbleClient { } = playObj; try { - await this.api.submitListen(playObj, { log: isDebugMode()}); + const result = await this.api.submitListen(playObj, { log: isDebugMode()}); if (newFromSource) { this.logger.info(`Scrobbled (New) => (${source}) ${buildTrackString(playObj)}`); } else { this.logger.info(`Scrobbled (Backlog) => (${source}) ${buildTrackString(playObj)}`); } - return playObj; + 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'}); - throw new UpstreamError(`Error occurred while making Koito API scrobble request: ${e.message}`, {cause: e, showStopper: !(e instanceof UpstreamError)}); + throw e; } } @@ -99,7 +99,7 @@ export default class KoitoScrobbler extends AbstractScrobbleClient { try { await this.api.submitListen(data, { listenType: 'playing_now'}); } catch (e) { - throw new UpstreamError(`Error occurred while making Koito API Playing Now request: ${e.message}`, {cause: e, showStopper: !(e instanceof UpstreamError)}); + throw e; } } } diff --git a/src/backend/scrobblers/LastfmScrobbler.ts b/src/backend/scrobblers/LastfmScrobbler.ts index bc08f597..710c3f03 100644 --- a/src/backend/scrobblers/LastfmScrobbler.ts +++ b/src/backend/scrobblers/LastfmScrobbler.ts @@ -9,6 +9,7 @@ import { LastfmClientConfig } from "../common/infrastructure/config/client/lastf import LastfmApiClient, { LastFMIgnoredScrobble, playToClientPayload, formatPlayObj, LASTFM_HOST, LASTFM_PATH } from "../common/vendor/LastfmApiClient.js"; import { Notifiers } from "../notifier/Notifiers.js"; import AbstractScrobbleClient, { nowPlayingUpdateByPlayDuration } from "./AbstractScrobbleClient.js"; +import { findCauseByReference } from "../utils/ErrorUtils.js"; export default class LastfmScrobbler extends AbstractScrobbleClient { @@ -95,17 +96,14 @@ export default class LastfmScrobbler extends AbstractScrobbleClient { } return respPlay; } catch (e) { - if(e instanceof LastFMIgnoredScrobble) { + const ignored = findCauseByReference(e, LastFMIgnoredScrobble); + if(ignored !== undefined) { await this.notifier.notify({title: `Client - ${capitalize(this.type)} - ${this.name} - Scrobble Ignored`, message: `Failed to scrobble => ${buildTrackString(playObj)} | ${e.message}`, priority: 'warn'}); } else { await this.notifier.notify({title: `Client - ${capitalize(this.type)} - ${this.name} - Scrobble Error`, message: `Failed to scrobble => ${buildTrackString(playObj)} | Error: ${e.message}`, priority: 'error'}); } this.logger.error({playInfo: buildTrackString(playObj), payload: playToClientPayload(playObj)}, `Scrobble Error (${sType})`); - if(!(e instanceof UpstreamError)) { - throw new UpstreamError(`Error received from ${this.upstreamType} API`, {cause: e, showStopper: true}); - } else { - throw e; - } + throw e; } } diff --git a/src/backend/scrobblers/ListenbrainzScrobbler.ts b/src/backend/scrobblers/ListenbrainzScrobbler.ts index 4edf5318..136433c7 100644 --- a/src/backend/scrobblers/ListenbrainzScrobbler.ts +++ b/src/backend/scrobblers/ListenbrainzScrobbler.ts @@ -3,10 +3,10 @@ import EventEmitter from "events"; import { PlayObject } from "../../core/Atomic.js"; import { buildTrackString, capitalize } from "../../core/StringUtils.js"; import { isNodeNetworkException } from "../common/errors/NodeErrors.js"; -import { UpstreamError } from "../common/errors/UpstreamError.js"; +import { hasUpstreamError, UpstreamError } from "../common/errors/UpstreamError.js"; import { FormatPlayObjectOptions } from "../common/infrastructure/Atomic.js"; import { ListenBrainzClientConfig } from "../common/infrastructure/config/client/listenbrainz.js"; -import { ListenbrainzApiClient, playToListenPayload } from "../common/vendor/ListenbrainzApiClient.js"; +import { ListenbrainzApiClient, playToListenPayload, playToSubmitPayload } from "../common/vendor/ListenbrainzApiClient.js"; import { ListenPayload } from '../common/vendor/listenbrainz/interfaces.js'; import { Notifiers } from "../notifier/Notifiers.js"; @@ -84,17 +84,17 @@ export default class ListenbrainzScrobbler extends AbstractScrobbleClient { } = playObj; try { - await this.api.submitListen(playObj, { log: isDebugMode()}); + const result = await this.api.submitListen(playObj, { log: isDebugMode()}); if (newFromSource) { this.logger.info(`Scrobbled (New) => (${source}) ${buildTrackString(playObj)}`); } else { this.logger.info(`Scrobbled (Backlog) => (${source}) ${buildTrackString(playObj)}`); } - return playObj; + 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'}); - throw new UpstreamError(`Error occurred while making Listenbrainz API scrobble request: ${e.message}`, {cause: e, showStopper: !(e instanceof UpstreamError)}); + throw e; } } @@ -103,7 +103,7 @@ export default class ListenbrainzScrobbler extends AbstractScrobbleClient { try { await this.api.submitListen(data, { listenType: 'playing_now'}); } catch (e) { - throw new UpstreamError(`Error occurred while making Listenbrainz API Playing Now request: ${e.message}`, {cause: e, showStopper: !(e instanceof UpstreamError)}); + throw e; } } } diff --git a/src/backend/scrobblers/MalojaScrobbler.ts b/src/backend/scrobblers/MalojaScrobbler.ts index 5afbdcfc..a7739ba5 100644 --- a/src/backend/scrobblers/MalojaScrobbler.ts +++ b/src/backend/scrobblers/MalojaScrobbler.ts @@ -12,6 +12,7 @@ import { import { Notifiers } from "../notifier/Notifiers.js"; import AbstractScrobbleClient from "./AbstractScrobbleClient.js"; import { MalojaApiClient, formatPlayObj as formatMalojaScrobbleToPlay, playToScrobblePayload } from "../common/vendor/maloja/MalojaApiClient.js"; +import { ScrobbleSubmitError } from "../common/errors/MSErrors.js"; const feat = ["ft.", "ft", "feat.", "feat", "featuring", "Ft.", "Ft", "Feat.", "Feat", "Featuring"]; @@ -98,31 +99,18 @@ export default class MalojaScrobbler extends AbstractScrobbleClient { const scrobbleData = playToScrobblePayload(playObj); - let scrobbledPlay: PlayObject; - try { - const [scrobbleResp, respBody, warnStr] = await this.api.scrobble(playObj); - - let warning = ''; - if (scrobbleResp === undefined) { - warning = `WARNING: Maloja did not return track data in scrobble response! Maybe it didn't scrobble correctly??`; - scrobbledPlay = playObj; - } else { - scrobbledPlay = this.formatPlayObj(scrobbleResp) - } + const result = await this.api.scrobble(playObj); const scrobbleInfo = `Scrobbled (${newFromSource ? 'New' : 'Backlog'}) => (${source}) ${buildTrackString(playObj)}`; - if (warning !== '') { - this.logger.warn(`${scrobbleInfo} | ${warning}`); - this.logger.debug(`Response: ${this.logger.debug(JSON.stringify(respBody))}`); + if (result.warnings?.length > 0) { + this.logger.warn(`${scrobbleInfo} | ${result.warnings.join(' | ')}`); } else { this.logger.info(scrobbleInfo); } - return scrobbledPlay; + 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' }); throw e; - } finally { - this.logger.debug(scrobbleData, 'Raw Payload'); } } } \ No newline at end of file diff --git a/src/backend/scrobblers/RockskyScrobbler.ts b/src/backend/scrobblers/RockskyScrobbler.ts index eaafae6d..d95c086b 100644 --- a/src/backend/scrobblers/RockskyScrobbler.ts +++ b/src/backend/scrobblers/RockskyScrobbler.ts @@ -3,7 +3,7 @@ import EventEmitter from "events"; import { PlayObject } from "../../core/Atomic.js"; import { buildTrackString, capitalize } from "../../core/StringUtils.js"; import { isNodeNetworkException } from "../common/errors/NodeErrors.js"; -import { UpstreamError } from "../common/errors/UpstreamError.js"; +import { hasUpstreamError, UpstreamError } from "../common/errors/UpstreamError.js"; import { FormatPlayObjectOptions } from "../common/infrastructure/Atomic.js"; import { ListenBrainzClientConfig } from "../common/infrastructure/config/client/listenbrainz.js"; import { ListenbrainzApiClient, playToListenPayload } from "../common/vendor/ListenbrainzApiClient.js"; @@ -12,8 +12,9 @@ import { Notifiers } from "../notifier/Notifiers.js"; import AbstractScrobbleClient from "./AbstractScrobbleClient.js"; import { isDebugMode } from "../utils.js"; -import { RockSkyApiClient } from "../common/vendor/RockSkyApiClient.js"; +import { RockSkyApiClient, 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 { @@ -84,10 +85,10 @@ export default class RockskyScrobbler extends AbstractScrobbleClient { } = playObj; try { - const resp = await this.api.submitListen(playObj, { log: isDebugMode()}); + const result = await this.api.submitListen(playObj, { log: isDebugMode()}); - if((resp.payload?.ignored_listens ?? 0) > 0) { - throw new UpstreamError('Scrobble was successfully submitted but Rocksky ignored it', {showStopper: false}); + 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 (newFromSource) { @@ -95,10 +96,10 @@ export default class RockskyScrobbler extends AbstractScrobbleClient { } else { this.logger.info(`Scrobbled (Backlog) => (${source}) ${buildTrackString(playObj)}`); } - return playObj; + 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'}); - throw new UpstreamError(`Error occurred while making Rocksky API scrobble request: ${e.message}`, {cause: e, showStopper: !(e instanceof UpstreamError)}); + throw e; } } @@ -106,7 +107,7 @@ export default class RockskyScrobbler extends AbstractScrobbleClient { try { await this.api.submitListen(data, { listenType: 'playing_now'}); } catch (e) { - throw new UpstreamError(`Error occurred while making Rocksky API Playing Now request: ${e.message}`, {cause: e, showStopper: !(e instanceof UpstreamError)}); + throw e; } } } diff --git a/src/backend/scrobblers/TealfmScrobbler.ts b/src/backend/scrobblers/TealfmScrobbler.ts index 81ced851..78b2744b 100644 --- a/src/backend/scrobblers/TealfmScrobbler.ts +++ b/src/backend/scrobblers/TealfmScrobbler.ts @@ -3,7 +3,7 @@ import EventEmitter from "events"; import { PlayObject } from "../../core/Atomic.js"; import { buildTrackString, capitalize } from "../../core/StringUtils.js"; import { isNodeNetworkException } from "../common/errors/NodeErrors.js"; -import { UpstreamError } from "../common/errors/UpstreamError.js"; +import { hasUpstreamError, UpstreamError } from "../common/errors/UpstreamError.js"; import { FormatPlayObjectOptions } from "../common/infrastructure/Atomic.js"; import { playToListenPayload } from "../common/vendor/ListenbrainzApiClient.js"; import { Notifiers } from "../notifier/Notifiers.js"; @@ -113,16 +113,16 @@ export default class TealScrobbler extends AbstractScrobbleClient { } = playObj; try { - await this.client.createScrobbleRecord(playToRecord(playObj)) + const res = await this.client.createScrobbleRecord(playToRecord(playObj)) if (newFromSource) { this.logger.info(`Scrobbled (New) => (${source}) ${buildTrackString(playObj)}`); } else { this.logger.info(`Scrobbled (Backlog) => (${source}) ${buildTrackString(playObj)}`); } - return playObj; + return res; } 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'}); - throw new UpstreamError(`Error occurred while making Teal API scrobble request: ${e.message}`, {cause: e, showStopper: !(e instanceof UpstreamError)}); + throw e; } } } diff --git a/src/backend/sources/SpotifySource.ts b/src/backend/sources/SpotifySource.ts index d92ea998..588ebe60 100644 --- a/src/backend/sources/SpotifySource.ts +++ b/src/backend/sources/SpotifySource.ts @@ -405,7 +405,7 @@ export default class SpotifySource extends MemoryPositionalSource { getPlayHistory = async (options: RecentlyPlayedOptions = {}) => { const {limit = 20} = options; const func = (api: SpotifyWebApi) => api.getMyRecentlyPlayedTracks({ - limit + limit: 1 }); const result = await this.callApi>(func); return result.body.items.map((x: PlayHistoryObject) => SpotifySource.formatPlayObj(x)).sort(sortByOldestPlayDate); diff --git a/src/backend/tests/scrobbler/TestScrobbler.ts b/src/backend/tests/scrobbler/TestScrobbler.ts index 19df2627..d13b9905 100644 --- a/src/backend/tests/scrobbler/TestScrobbler.ts +++ b/src/backend/tests/scrobbler/TestScrobbler.ts @@ -21,8 +21,8 @@ export class TestScrobbler extends AbstractScrobbleClient { return this.testRecentScrobbles; } - doScrobble(playObj: PlayObject): Promise { - return Promise.resolve(playObj); + doScrobble(playObj: PlayObject) { + return Promise.resolve({payload: {}, mergedScrobble: playObj}); } alreadyScrobbled = async (playObj: PlayObject, log?: boolean): Promise => { diff --git a/src/core/Atomic.ts b/src/core/Atomic.ts index 511ab119..6e117ad3 100644 --- a/src/core/Atomic.ts +++ b/src/core/Atomic.ts @@ -287,7 +287,13 @@ export interface PlayLifecycle { input?: object original: PlayObjectLifecycleless steps: LifecycleStep[] - scrobble?: object + scrobble?: { + payload?: ScrobblePayload + warnings?: string[] + error?: Error + response?: ScrobbleResponse + mergedScrobble?: PlayObjectLifecycleless + } } export interface LifecycleStep { @@ -297,6 +303,16 @@ export interface LifecycleStep { inputs?: LifecycleInput[] } +export type ScrobblePayload = object | string; +export type ScrobbleResponse = object | string; + +export interface ScrobbleActionResult { + payload: ScrobblePayload, + response?: ScrobbleResponse, + mergedScrobble?: PlayObject + warnings?: string[] +} + export type ScrobbleTsSOC = 1 | 2; export const SCROBBLE_TS_SOC_START: ScrobbleTsSOC = 1; -- 2.51.2