From a41f4c56f73dd31c7b8aa1487b30c7a00837d8aa Mon Sep 17 00:00:00 2001 From: FoxxMD Date: Fri, 6 Mar 2026 16:57:09 +0000 Subject: [PATCH] feat(scrobbler): Implement play-dependent, arbitrary timerange fetching for duplicate checks * Remove design bias towards "recent" scrobbles with limited timerange and complex shouldRefresh logic * Replace with logic to fetch scrobbles from timeranges that are relevant to queued/dead scrobbles * timeranges are grouped by closely timestamped queued scrobbles * timeranges are cached based on staleAfter --- package-lock.json | 88 +++++++ package.json | 2 + src/backend/common/Cache.ts | 2 + src/backend/common/infrastructure/Atomic.ts | 13 +- .../scrobblers/AbstractScrobbleClient.ts | 238 +++++------------ src/backend/scrobblers/DiscordScrobbler.ts | 8 +- src/backend/scrobblers/KoitoScrobbler.ts | 10 +- src/backend/scrobblers/LastfmScrobbler.ts | 10 - .../scrobblers/ListenbrainzScrobbler.ts | 10 +- src/backend/scrobblers/MalojaScrobbler.ts | 10 +- src/backend/scrobblers/RockskyScrobbler.ts | 4 +- src/backend/scrobblers/TealfmScrobbler.ts | 4 +- src/backend/tests/scrobbler/TestScrobbler.ts | 9 +- .../tests/scrobbler/scrobblers.test.ts | 239 +++++++++++------- src/backend/utils/DataUtils.ts | 11 + src/backend/utils/ListenFetchUtils.ts | 71 +++++- src/backend/utils/PlayComparisonUtils.ts | 5 + 17 files changed, 413 insertions(+), 321 deletions(-) diff --git a/package-lock.json b/package-lock.json index e4775fd8..7f67f58d 100644 --- a/package-lock.json +++ b/package-lock.json @@ -140,6 +140,7 @@ "@types/react": "^18.2.18", "@types/react-dom": "^18.2.7", "@types/react-window": "^1.8.5", + "@types/sinon": "^21.0.0", "@types/spotify-web-api-node": "^5.0.7", "@types/superagent": "^8.1.9", "@types/xml2js": "^0.4.11", @@ -153,6 +154,7 @@ "mockdate": "^3.0.5", "msw": "^2.1.2", "nodemon": "^3.0.3", + "sinon": "^21.0.2", "ts-essentials": "^9.1.2", "typescript": "5.5.4", "typescript-eslint": "^7.0.1", @@ -3039,6 +3041,47 @@ "url": "https://github.com/sponsors/sindresorhus" } }, + "node_modules/@sinonjs/commons": { + "version": "3.0.1", + "resolved": "https://registry.npmjs.org/@sinonjs/commons/-/commons-3.0.1.tgz", + "integrity": "sha512-K3mCHKQ9sVh8o1C9cxkwxaOmXoAMlDxC1mYyHrjqOWEcBjYr76t96zL2zlj5dUGZ3HSw240X1qgH3Mjf1yJWpQ==", + "dev": true, + "license": "BSD-3-Clause", + "dependencies": { + "type-detect": "4.0.8" + } + }, + "node_modules/@sinonjs/commons/node_modules/type-detect": { + "version": "4.0.8", + "resolved": "https://registry.npmjs.org/type-detect/-/type-detect-4.0.8.tgz", + "integrity": "sha512-0fr/mIH1dlO+x7TlcMy+bIDqKPsw/70tVyeHW787goQjhmqaZe10uwLujubK9q9Lg6Fiho1KUKDYz0Z7k7g5/g==", + "dev": true, + "license": "MIT", + "engines": { + "node": ">=4" + } + }, + "node_modules/@sinonjs/fake-timers": { + "version": "15.1.1", + "resolved": "https://registry.npmjs.org/@sinonjs/fake-timers/-/fake-timers-15.1.1.tgz", + "integrity": "sha512-cO5W33JgAPbOh07tvZjUOJ7oWhtaqGHiZw+11DPbyqh2kHTBc3eF/CjJDeQ4205RLQsX6rxCuYOroFQwl7JDRw==", + "dev": true, + "license": "BSD-3-Clause", + "dependencies": { + "@sinonjs/commons": "^3.0.1" + } + }, + "node_modules/@sinonjs/samsam": { + "version": "9.0.2", + "resolved": "https://registry.npmjs.org/@sinonjs/samsam/-/samsam-9.0.2.tgz", + "integrity": "sha512-H/JSxa4GNKZuuU41E3b8Y3tbSEx8y4uq4UH1C56ONQac16HblReJomIvv3Ud7ANQHQmkeSowY49Ij972e/pGxQ==", + "dev": true, + "license": "BSD-3-Clause", + "dependencies": { + "@sinonjs/commons": "^3.0.1", + "type-detect": "^4.1.0" + } + }, "node_modules/@standard-schema/spec": { "version": "1.0.0", "resolved": "https://registry.npmjs.org/@standard-schema/spec/-/spec-1.0.0.tgz", @@ -4066,6 +4109,23 @@ "@types/send": "*" } }, + "node_modules/@types/sinon": { + "version": "21.0.0", + "resolved": "https://registry.npmjs.org/@types/sinon/-/sinon-21.0.0.tgz", + "integrity": "sha512-+oHKZ0lTI+WVLxx1IbJDNmReQaIsQJjN2e7UUrJHEeByG7bFeKJYsv1E75JxTQ9QKJDp21bAa/0W2Xo4srsDnw==", + "dev": true, + "license": "MIT", + "dependencies": { + "@types/sinonjs__fake-timers": "*" + } + }, + "node_modules/@types/sinonjs__fake-timers": { + "version": "15.0.1", + "resolved": "https://registry.npmjs.org/@types/sinonjs__fake-timers/-/sinonjs__fake-timers-15.0.1.tgz", + "integrity": "sha512-Ko2tjWJq8oozHzHV+reuvS5KYIRAokHnGbDwGh/J64LntgpbuylF74ipEL24HCyRjf9FOlBiBHWBR1RlVKsI1w==", + "dev": true, + "license": "MIT" + }, "node_modules/@types/spotify-api": { "version": "0.0.25", "resolved": "https://registry.npmjs.org/@types/spotify-api/-/spotify-api-0.0.25.tgz", @@ -12282,6 +12342,34 @@ "node": ">=10" } }, + "node_modules/sinon": { + "version": "21.0.2", + "resolved": "https://registry.npmjs.org/sinon/-/sinon-21.0.2.tgz", + "integrity": "sha512-VHV4UaoxIe5jrMd89Y9duI76T5g3Lp+ET+ctLhLDaZtSznDPah1KKpRElbdBV4RwqWSw2vadFiVs9Del7MbVeQ==", + "dev": true, + "license": "BSD-3-Clause", + "dependencies": { + "@sinonjs/commons": "^3.0.1", + "@sinonjs/fake-timers": "^15.1.1", + "@sinonjs/samsam": "^9.0.2", + "diff": "^8.0.3", + "supports-color": "^7.2.0" + }, + "funding": { + "type": "opencollective", + "url": "https://opencollective.com/sinon" + } + }, + "node_modules/sinon/node_modules/diff": { + "version": "8.0.3", + "resolved": "https://registry.npmjs.org/diff/-/diff-8.0.3.tgz", + "integrity": "sha512-qejHi7bcSD4hQAZE0tNAawRK1ZtafHDmMTMkrrIGgSLl7hTnQHmKCeB45xAcbfTqK2zowkM3j3bHt/4b/ARbYQ==", + "dev": true, + "license": "BSD-3-Clause", + "engines": { + "node": ">=0.3.1" + } + }, "node_modules/slash": { "version": "2.0.0", "resolved": "https://registry.npmjs.org/slash/-/slash-2.0.0.tgz", diff --git a/package.json b/package.json index d0f627e0..261afe4b 100644 --- a/package.json +++ b/package.json @@ -175,6 +175,7 @@ "@types/react": "^18.2.18", "@types/react-dom": "^18.2.7", "@types/react-window": "^1.8.5", + "@types/sinon": "^21.0.0", "@types/spotify-web-api-node": "^5.0.7", "@types/superagent": "^8.1.9", "@types/xml2js": "^0.4.11", @@ -188,6 +189,7 @@ "mockdate": "^3.0.5", "msw": "^2.1.2", "nodemon": "^3.0.3", + "sinon": "^21.0.2", "ts-essentials": "^9.1.2", "typescript": "5.5.4", "typescript-eslint": "^7.0.1", diff --git a/src/backend/common/Cache.ts b/src/backend/common/Cache.ts index 329ede20..c61d11d8 100644 --- a/src/backend/common/Cache.ts +++ b/src/backend/common/Cache.ts @@ -50,6 +50,7 @@ export class MSCache { cacheAuth: Cacheable; regexCache: ReturnType; cacheTransform: Cacheable; + cacheClientScrobbles: Cacheable; logger: Logger; @@ -96,6 +97,7 @@ export class MSCache { this.regexCache = cacheFunctions(this.config.regex); this.cacheTransform = new Cacheable({primary: initMemoryCache({lruSize: 500})}); + this.cacheClientScrobbles = new Cacheable({primary: initMemoryCache({lruSize: 100, ttl: '5m'})}); } init = async () => { diff --git a/src/backend/common/infrastructure/Atomic.ts b/src/backend/common/infrastructure/Atomic.ts index b747ef87..5e27f2df 100644 --- a/src/backend/common/infrastructure/Atomic.ts +++ b/src/backend/common/infrastructure/Atomic.ts @@ -347,6 +347,10 @@ export interface Authenticatable { testAuth: () => Promise } +export interface ScrobbleRangeFetchable { + getScrobblesForTimeRange: TimeRangeListensFetcher +} + export interface MdnsDeviceInfo { name: string type: string @@ -501,4 +505,11 @@ export const hasPagelessTimeRangeListens = (obj: Object): obj is PagelessTimeRan export type PaginatedTimeRangeSource = PaginatedTimeRangeListens | PagelessTimeRangeListens; export type PaginatedSource = PaginatedListens | PaginatedTimeRangeSource; -export type TimeRangeListensFetcher = (opts: PaginatedTimeRangeCommonOptions | PaginatedListensTimeRangeOptions) => Promise \ No newline at end of file +export type TimeRangeListensFetcher = (opts: PaginatedTimeRangeCommonOptions | PaginatedListensTimeRangeOptions) => Promise + +export interface ScrobbleRangeResult { + plays: PlayObject[] + fetchedAt: Dayjs +} + +export const REFRESH_STALE_DEFAULT = 60; \ No newline at end of file diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index 502c7a89..19a61752 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -9,13 +9,13 @@ import { NowPlayingUpdateThreshold, PlayObject, PlayObjectLifecycleless, - QueuedScrobble, ScrobbleActionResult, PlayMatchResult, ScrobblePayload, ScrobbleResponse, SourcePlayerObj, TA_DURING, + QueuedScrobble, ScrobbleActionResult, PlayMatchResult, SourcePlayerObj, TA_DURING, TA_FUZZY, TrackStringOptions } from "../../core/Atomic.js"; import { buildTrackString, capitalize, truncateStringToLength } from "../../core/StringUtils.js"; import AbstractComponent from "../common/AbstractComponent.js"; -import { hasUpstreamError, UpstreamError } from "../common/errors/UpstreamError.js"; +import { hasUpstreamError } from "../common/errors/UpstreamError.js"; import { ARTIST_WEIGHT, Authenticatable, @@ -24,10 +24,12 @@ import { DEFAULT_RETRY_MULTIPLIER, DUP_SCORE_THRESHOLD, FormatPlayObjectOptions, - PlayPlatformId, + PaginatedTimeRangeOptions, + REFRESH_STALE_DEFAULT, ScrobbledPlayObject, SourceIdentifier, TIME_WEIGHT, + TimeRangeListensFetcher, TITLE_WEIGHT, } from "../common/infrastructure/Atomic.js"; import { CommonClientConfig, NowPlayingOptions, UpstreamRefreshOptions } from "../common/infrastructure/config/client/index.js"; @@ -35,14 +37,10 @@ import { TRANSFORM_HOOK } from "../common/infrastructure/Transform.js"; import { Notifiers } from "../notifier/Notifiers.js"; import { comparingMultipleArtists, - genGroupId, - genGroupIdStr, - genGroupIdStrFromPlay, isDebugMode, parseBool, playObjDataMatch, pollingBackoff, - setIntersection, sleep, sortByOldestPlayDate, } from "../utils.js"; @@ -55,7 +53,6 @@ import { } from "../utils/TimeUtils.js"; import { WebhookPayload } from "../common/infrastructure/config/health/webhooks.js"; import { AsyncTask, SimpleIntervalJob, Task, ToadScheduler } from "toad-scheduler"; -import { MSCache } from "../common/Cache.js"; import { getRoot } from "../ioc.js"; import { rehydratePlay } from "../utils/CacheUtils.js"; import { findAsyncSequential, staggerMapper } from "../utils/AsyncUtils.js"; @@ -65,8 +62,7 @@ import { normalizeStr } from "../utils/StringUtils.js"; import prom, { Counter, Gauge } from 'prom-client'; import { ScrobbleSubmitError } from "../common/errors/MSErrors.js"; import {serializeError} from 'serialize-error'; -import { redactString } from "@foxxmd/redact-string"; -import clone from "clone"; +import { DEFAULT_NEW_PADDING, groupPlaysToTimeRanges } from "../utils/ListenFetchUtils.js"; type PlatformMappedPlays = Map; type NowPlayingQueue = Map; @@ -83,14 +79,10 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i protected MAX_STORED_SCROBBLES = 40; protected MAX_INITIAL_SCROBBLES_FETCH = this.MAX_STORED_SCROBBLES; - #recentScrobblesList: PlayObject[] = []; + scrobbleSOTRanges: PaginatedTimeRangeOptions[] = []; scrobbledPlayObjs: FixedSizeList; - lastScrobbledPlayDate?: Dayjs; - newestScrobbleTime?: Dayjs - oldestScrobbleTime?: Dayjs tracksScrobbled: number = 0; - lastScrobbleCheck: Dayjs = dayjs(0) lastScrobbleAttempt: Dayjs = dayjs(0) upstreamRefresh: MarkOptional, 'refreshInitialCount'>; checkExistingScrobbles: boolean; @@ -142,7 +134,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i refreshEnabled = true, refreshInitialCount, refreshMinInterval = 5, - refreshStaleAfter = 60, + refreshStaleAfter = REFRESH_STALE_DEFAULT, checkExistingScrobbles = true, verbose = {}, } = {}, @@ -185,16 +177,6 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i this.scrobbledCounter = clientMetrics.scrobbled; } - set recentScrobbles(scrobbles: PlayObject[]) { - const sorted = [...scrobbles]; - sorted.sort(sortByOldestPlayDate); - this.#recentScrobblesList = sorted; - } - - get recentScrobbles() { - return this.#recentScrobblesList; - } - protected getIdentifier() { return `${capitalize(this.type)} - ${this.name}` } @@ -357,109 +339,32 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i } = this.config; this.initializeNowPlaying(); - - let initialLimit = refreshInitialCount; - if(refreshInitialCount > this.MAX_INITIAL_SCROBBLES_FETCH) { - this.logger.warn(`Defined initial scrobbles count (${refreshInitialCount}) higher than maximum allowed (${this.MAX_INITIAL_SCROBBLES_FETCH}). Will use max instead.`); - initialLimit = this.MAX_INITIAL_SCROBBLES_FETCH; - } - - this.logger.verbose(`Fetching up to ${initialLimit} initial scrobbles...`); - await this.refreshScrobbles(initialLimit); - this.lastScrobbledPlayDate = this.newestScrobbleTime; } - refreshScrobbles = async (limit: number = this.MAX_STORED_SCROBBLES) => { - if (this.upstreamRefresh.refreshEnabled) { - this.logger.debug('Refreshing recent scrobbles'); - const recent = await this.getScrobblesForRefresh(limit); - this.logger.debug(`Found ${recent.length} recent scrobbles`); - this.recentScrobbles = recent; - if (this.recentScrobbles.length > 0) { - const [{data: {playDate: newestScrobbleTime = dayjs()} = {}} = {}] = this.recentScrobbles.slice(-1); - const [{data: {playDate: oldestScrobbleTime = dayjs()} = {}} = {}] = this.recentScrobbles.slice(0, 1); - this.newestScrobbleTime = newestScrobbleTime; - this.oldestScrobbleTime = oldestScrobbleTime; - - this.filterScrobbledTracks(); - } - } - this.lastScrobbleCheck = dayjs(); - } - - protected abstract getScrobblesForRefresh(limit: number): Promise; - - shouldRefreshScrobble = () => { - const { - refreshStaleAfter, - refreshMinInterval, - refreshEnabled - } = this.upstreamRefresh; - - if (!refreshEnabled) { - this.logger.debug({labels: ['Upstream Refresh']}, `Should NOT refresh => refreshEnabled is false`); - return false; - } - - if(this.queuedScrobbles.length === 0) { - this.logger.debug({labels: ['Upstream Refresh']}, `Should NOT refresh => no scrobbles in queue!`); - return false; - } - - const queuedPlayedDate = this.getLatestQueuePlayDate(); - - // if newest queued play was played more recently than the last time we refreshed upstream scrobbles - if (this.scrobblesLastCheckedAt().unix() < queuedPlayedDate.unix()) { - if(!this.scrobblesRefreshMinIntervalPassed()) { - this.logger.debug({labels: ['Upstream Refresh']}, `Should refresh but WILL NOT => queued scrobble playDate is newer than last refresh but refreshMinInterval (${refreshMinInterval}ms) has not passed since last check`); - return false; - } else { - this.logger.debug({labels: ['Upstream Refresh']}, 'Should refresh => newest queued scrobble playDate is newer than last refresh'); - return true; - } - } - - // if the play date of the last Play scrobbled is *newer* - // than the queued scrobble we are about to scrobble - // then we are inserting a scrobble out of order which can happen if - // * backlogging and upstream returned plays out of order - // * processing dead letter queue - // * two sources have different history - // - // in all cases we probably want to refresh - if(this.lastScrobbledPlayDate !== undefined && this.queuedScrobbles[0].play.meta.newFromSource && this.lastScrobbledPlayDate.unix() > this.queuedScrobbles[0].play.data.playDate.unix()) { - if(!this.scrobblesRefreshMinIntervalPassed()) { - this.logger.debug({labels: ['Upstream Refresh']}, `Should refresh but WILL NOT => queued scrobble playDate is older than last scrobbled play (out-of-order insert) but refreshMinInterval (${refreshMinInterval}ms) has not passed since last check`); - return false; - } else { - this.logger.debug({labels: ['Upstream Refresh']}, 'Should refresh => queued scrobble playDate is older than last scrobbled play (out-of-order insert)'); - return true; - } - } - - // if it's been X seconds since we last refreshed - if(refreshStaleAfter !== undefined) { - const diff = dayjs().diff(this.scrobblesLastCheckedAt(), 's'); - if(diff > refreshStaleAfter) { - this.logger.debug({labels: ['Upstream Refresh']}, `Should refresh => last refresh (${diff}s ago) was longer than refreshStaleAfter (${refreshStaleAfter}s)`); - return true; - } - } + abstract getScrobblesForTimeRange: TimeRangeListensFetcher; - return false; + handleQueuedScrobbleRanges = () => { + this.scrobbleSOTRanges = groupPlaysToTimeRanges(this.queuedScrobbles.map(x => x.play).concat(this.deadLetterScrobbles.map(x => x.play)), this.scrobbleSOTRanges, {staleNowBuffer: this.config.options?.refreshStaleAfter}); } - protected scrobblesLastCheckedAt = () => this.lastScrobbleCheck - protected scrobblesLastCheckedAtDiff = () => dayjs().diff(this.scrobblesLastCheckedAt(), 'ms') - protected scrobblesRefreshMinIntervalPassed = () => { - const { - refreshMinInterval, - } = this.upstreamRefresh; - return this.scrobblesLastCheckedAtDiff() >= refreshMinInterval; + getSOTScrobblesForPlay = async (play: PlayObject): Promise => { + let range: PaginatedTimeRangeOptions = this.scrobbleSOTRanges.find(x => x.from <= play.data.playDate.unix() && x.to > Math.min(dayjs().subtract(this.config.options?.refreshStaleAfter ?? REFRESH_STALE_DEFAULT, 's').unix(), play.data.playDate.unix())); + if(range === undefined) { + this.logger.warn(`No Scrobble SOT range found! Should have been handled before this. Creating a new one for ${buildTrackString(play)}`); + range = { + from: play.data.playDate.subtract(DEFAULT_NEW_PADDING).unix(), + to: Math.min(play.data.playDate.add(DEFAULT_NEW_PADDING).unix(), dayjs().subtract(this.config.options?.refreshStaleAfter ?? REFRESH_STALE_DEFAULT, 's').unix()) + }; + this.scrobbleSOTRanges.push(range); + } + const cacheKey = `${this.name}-scrobbleRange-${range.from}-${range.to}`; + const plays = this.cache.cacheClientScrobbles.getOrSet(cacheKey, async () => { + return await this.getScrobblesForTimeRange(range); + }, {ttl: this.config.options?.refreshStaleAfter ?? REFRESH_STALE_DEFAULT}); + return plays; } - public async alreadyScrobbled(playObj: PlayObject, log?: boolean): Promise<[boolean, PlayMatchResult]> { - const result = await this.existingScrobble(playObj); + const result = await this.existingScrobble(playObj, await this.getSOTScrobblesForPlay(playObj)); return [result.match, result]; } @@ -468,39 +373,13 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i return obj; } - // time frame is valid as long as the play date for the source track is newer than the oldest play time from the scrobble client - // ...this is assuming the scrobble client is returning "most recent" scrobbles - timeFrameIsValid = (playObj: PlayObject) => { - - if(this.oldestScrobbleTime === undefined) { - return [true, '']; - } - - const { - data: { - playDate, - } = {}, - } = playObj; - const validTime = playDate.isAfter(this.oldestScrobbleTime); - let log = ''; - if (!validTime) { - const dur = dayjs.duration(Math.abs(playDate.diff(this.oldestScrobbleTime))).humanize(false); - log = `occurred ${dur} before the oldest scrobble returned by this client (${this.oldestScrobbleTime.format()})`; - } - return [validTime, log] - } - addScrobbledTrack = (playObj: PlayObject, scrobbledPlay: PlayObjectLifecycleless) => { this.scrobbledPlayObjs.add({play: playObj, scrobble: scrobbledPlay}); this.scrobbledCounter.labels(this.getPrometheusLabels()).inc(); - this.lastScrobbledPlayDate = playObj.data.playDate; + //this.lastScrobbledPlayDate = playObj.data.playDate; this.tracksScrobbled++; } - filterScrobbledTracks = () => { - this.scrobbledPlayObjs = new FixedSizeList(this.MAX_STORED_SCROBBLES, this.scrobbledPlayObjs.data.filter(x => this.timeFrameIsValid(x.play)[0])) ; - } - getScrobbledPlays = () => this.scrobbledPlayObjs.data.map(x => x.play) findExistingSubmittedPlayObj = async (playObjPre: PlayObject): Promise<([undefined, undefined] | [ScrobbledPlayObject, ScrobbledPlayObject[]])> => { @@ -523,7 +402,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i return [matchPlayDate, dtInvariantMatches]; } - existingScrobble = async (playObjPre: PlayObject): Promise => { + existingScrobble = async (playObjPre: PlayObject, existingScrobbles: PlayObject[]): Promise => { const result: PlayMatchResult = { match: false, @@ -568,8 +447,8 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i // if no recent scrobbles found then assume we haven't submitted it // (either user doesnt want to check history or there is no history to check!) - if (this.recentScrobbles.length === 0) { - this.dupeLogger.trace(`${buildTrackString(playObj, scoreTrackOpts)} => No Match because no recent scrobbles returned from API`); + if (existingScrobbles.length === 0) { + this.dupeLogger.trace(`${buildTrackString(playObj, scoreTrackOpts)} => No Match because no existing scrobbles returned from API`); result.reason = 'no recent scrobbles returned from API'; return result; } @@ -584,7 +463,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i const looseTimeAccuracy = playObj.data.repeat ? [TA_DURING] : [TA_FUZZY, TA_DURING]; - existingScrobble = await findAsyncSequential(this.recentScrobbles, async (xPre) => { + existingScrobble = await findAsyncSequential(existingScrobbles, async (xPre) => { const x = await this.transformPlay(xPre, TRANSFORM_HOOK.existing); @@ -809,16 +688,21 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i try { this.scrobbling = true; + if(!this.upstreamRefresh.refreshEnabled) { + this.logger.verbose('Scrobble refresh is DISABLED. All queued scrobbles will likely always be scrobbled (nothing to check duplicates against).'); + } while (!this.shouldStopScrobbleProcessing()) { + let queueEmpty = this.queuedScrobbles.length === 0; while (this.queuedScrobbles.length > 0) { - if (this.shouldRefreshScrobble()) { - await this.refreshScrobbles(); + this.handleQueuedScrobbleRanges(); + if(!this.upstreamRefresh.refreshEnabled) { + this.logger.trace('Scrobble refresh is DISABLED.'); } + const currQueuedPlay = this.queuedScrobbles.shift(); - const [timeFrameValid, timeFrameValidLog] = this.timeFrameIsValid(currQueuedPlay.play); - if (timeFrameValid) { - const [matched, matchResult] = await this.alreadyScrobbled(currQueuedPlay.play); + //const [timeFrameValid, timeFrameValidLog] = this.timeFrameIsValid(currQueuedPlay.play); + const matchResult = await this.existingScrobble(currQueuedPlay.play, !this.upstreamRefresh.refreshEnabled ? [] : await this.getSOTScrobblesForPlay(currQueuedPlay.play)); const { scrobble = {}, ...lifeRest @@ -830,7 +714,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i match: matchResult } } - if(!matched) { + if(!matchResult.match) { const transformedScrobble = await this.transformPlay(currQueuedPlay.play, TRANSFORM_HOOK.postCompare); if(transformedScrobble.meta.lifecycle === undefined) { transformedScrobble.meta.lifecycle = { @@ -866,13 +750,13 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i } } } - } else if (!timeFrameValid) { - this.logger.debug(`Will not scrobble ${buildTrackString(currQueuedPlay.play)} from Source '${currQueuedPlay.source}' because it ${timeFrameValidLog}`); - } this.updateQueuedScrobblesCache(); this.queuedGauge.labels(this.getPrometheusLabels()).set(this.queuedScrobbles.length); this.emitEvent('scrobbleDequeued', {queuedScrobble: currQueuedPlay}) } + if(!queueEmpty) { + this.emitEvent('queueEmptied', {}); + } await sleep(this.scrobbleSleep); } if (this.shouldStopScrobbleProcessing()) { @@ -905,10 +789,14 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i const processable = this.deadLetterScrobbles.filter(x => x.retries < retries); const queueStatus = `${processable.length} of ${this.deadLetterScrobbles.length} dead scrobbles have less than ${retries} retries, ${processable.length === 0 ? 'will skip processing.': 'processing now...'}`; if (processable.length === 0) { - this.logger.verbose(queueStatus, {leaf: 'Dead Letter'}); + this.logger.verbose({labels: 'Dead Letter'}, queueStatus); return; } - this.logger.info(queueStatus, {leaf: 'Dead Letter'}); + this.logger.info({labels: 'Dead Letter'}, queueStatus); + if(!this.upstreamRefresh.refreshEnabled) { + this.logger.verbose({labels: 'Dead Letter'}, 'Scrobble refresh is DISABLED. All dead scrobbles will likely always be scrobbled (nothing to check duplicates against).'); + } + this.handleQueuedScrobbleRanges(); const removedIds = []; for (const deadScrobble of this.deadLetterScrobbles) { @@ -920,7 +808,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i } } if (removedIds.length > 0) { - this.logger.info(`Removed ${removedIds.length} scrobbles from dead letter queue`, {leaf: 'Dead Letter'}); + this.logger.info({labels: 'Dead Letter'}, `Removed ${removedIds.length} scrobbles from dead letter queue`); } } @@ -929,15 +817,14 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i const deadScrobble = this.deadLetterScrobbles[deadScrobbleIndex]; if (!(await this.isReady())) { - this.logger.warn('Cannot process dead letter scrobble because client is not ready.', {leaf: 'Dead Letter'}); + this.logger.warn({labels: 'Dead Letter'}, 'Cannot process dead letter scrobble because client is not ready.'); return [false, deadScrobble]; } - if (this.getLatestQueuePlayDate() !== undefined && this.scrobblesLastCheckedAt().unix() < this.getLatestQueuePlayDate().unix()) { - await this.refreshScrobbles(); - } - const [timeFrameValid, timeFrameValidLog] = this.timeFrameIsValid(deadScrobble.play); - if (timeFrameValid) { - const [matched, matchResult] = await this.alreadyScrobbled(deadScrobble.play); + // if (this.getLatestQueuePlayDate() !== undefined && this.scrobblesLastCheckedAt().unix() < this.getLatestQueuePlayDate().unix()) { + // await this.refreshScrobbles(); + // } + //const [timeFrameValid, timeFrameValidLog] = this.timeFrameIsValid(deadScrobble.play); + const matchResult = await this.existingScrobble(deadScrobble.play, !this.upstreamRefresh.refreshEnabled ? [] : await this.getSOTScrobblesForPlay(deadScrobble.play)); const { scrobble = {}, ...lifeRest @@ -946,10 +833,10 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i ...lifeRest, scrobble: { ...scrobble, - match: matchResult[1] + match: matchResult } } - if(!matched) { + if(!matchResult.match) { const transformedScrobble = await this.transformPlay(deadScrobble.play, TRANSFORM_HOOK.postCompare); try { const scrobbledPlay = await this.scrobble(transformedScrobble); @@ -977,9 +864,6 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i await sleep(1000); } } - } else if (!timeFrameValid) { - this.logger.debug(`Will not scrobble ${buildTrackString(deadScrobble.play)} from Source '${deadScrobble.source}' because it ${timeFrameValidLog}`, {leaf: 'Dead Letter'}); - } if(deadScrobble !== undefined) { this.removeDeadLetterScrobble(deadScrobble.id) } diff --git a/src/backend/scrobblers/DiscordScrobbler.ts b/src/backend/scrobblers/DiscordScrobbler.ts index 6d7265d1..2d2c2c11 100644 --- a/src/backend/scrobblers/DiscordScrobbler.ts +++ b/src/backend/scrobblers/DiscordScrobbler.ts @@ -1,7 +1,7 @@ import { Logger } from "@foxxmd/logging"; import EventEmitter from "events"; import { PlayMatchResult, PlayObject, SourcePlayerObj } from "../../core/Atomic.js"; -import { CALCULATED_PLAYER_STATUSES, FormatPlayObjectOptions, REPORTED_PLAYER_STATUSES, ReportedPlayerStatus, SINGLE_USER_PLATFORM_ID_STR } from "../common/infrastructure/Atomic.js"; +import { CALCULATED_PLAYER_STATUSES, FormatPlayObjectOptions, REPORTED_PLAYER_STATUSES, ReportedPlayerStatus, SINGLE_USER_PLATFORM_ID_STR, TimeRangeListensFetcher } from "../common/infrastructure/Atomic.js"; import { Notifiers } from "../notifier/Notifiers.js"; import AbstractScrobbleClient, { nowPlayingUpdateByPlayDuration } from "./AbstractScrobbleClient.js"; @@ -30,6 +30,8 @@ export default class DiscordScrobbler extends AbstractScrobbleClient { this.nowPlayingMinThreshold = (_) => 5; } + getScrobblesForTimeRange = async (_) => []; + formatPlayObj = (obj: any, options: FormatPlayObjectOptions = {}) => obj; protected async doBuildInitData(): Promise { @@ -114,10 +116,6 @@ export default class DiscordScrobbler extends AbstractScrobbleClient { } } - getScrobblesForRefresh = async (limit: number) => { - return []; - } - queueScrobble = async (data: PlayObject | PlayObject[], source: string) => { // discord does not handle scrobbles, only Now Playing // so don't bother queueing any scrobbles as we don't want to cache them diff --git a/src/backend/scrobblers/KoitoScrobbler.ts b/src/backend/scrobblers/KoitoScrobbler.ts index d8ba226c..c055b191 100644 --- a/src/backend/scrobblers/KoitoScrobbler.ts +++ b/src/backend/scrobblers/KoitoScrobbler.ts @@ -8,7 +8,7 @@ import { FormatPlayObjectOptions, TimeRangeListensFetcher } from "../common/infr import { playToListenPayload } from "../common/vendor/ListenbrainzApiClient.js"; import { Notifiers } from "../notifier/Notifiers.js"; -import AbstractScrobbleClient, { shouldUpdatePlayingNowPlatformWhenPlayingOnly } from "./AbstractScrobbleClient.js"; +import AbstractScrobbleClient from "./AbstractScrobbleClient.js"; import { isDebugMode } from "../utils.js"; import { KoitoClientConfig } from "../common/infrastructure/config/client/koito.js"; import { KoitoApiClient, listenObjectResponseToPlay } from "../common/vendor/koito/KoitoApiClient.js"; @@ -69,14 +69,6 @@ export default class KoitoScrobbler extends AbstractScrobbleClient { } } - getScrobblesForRefresh = async (limit: number) => { - if(this.queuedScrobbles.length === 0) { - return await this.getScrobblesForTimeRange({limit, cursor: 0}); - } else { - return await this.getScrobblesForTimeRange({limit, cursor: 0, from: this.queuedScrobbles[0].play.data.playDate.unix(), to: dayjs().unix()}); - } - } - doScrobble = async (playObj: PlayObject) => { const { meta: { diff --git a/src/backend/scrobblers/LastfmScrobbler.ts b/src/backend/scrobblers/LastfmScrobbler.ts index 8712f70d..e4e1bdf5 100644 --- a/src/backend/scrobblers/LastfmScrobbler.ts +++ b/src/backend/scrobblers/LastfmScrobbler.ts @@ -60,16 +60,6 @@ export default class LastfmScrobbler extends AbstractScrobbleClient { } } - getScrobblesForRefresh = async (limit: number) => { - let plays: PlayObject[] = []; - if(this.queuedScrobbles.length === 0) { - plays = await this.getScrobblesForTimeRange({limit, cursor: 1}); - } else { - plays = await this.getScrobblesForTimeRange({limit, cursor: 1, from: this.queuedScrobbles[0].play.data.playDate.unix(), to: dayjs().unix()}); - } - return plays.filter(x => !x.meta.nowPlaying); - } - cleanSourceSearchTitle = (playObj: PlayObject) => { const { data: { diff --git a/src/backend/scrobblers/ListenbrainzScrobbler.ts b/src/backend/scrobblers/ListenbrainzScrobbler.ts index 03f25aee..011cc192 100644 --- a/src/backend/scrobblers/ListenbrainzScrobbler.ts +++ b/src/backend/scrobblers/ListenbrainzScrobbler.ts @@ -69,15 +69,7 @@ export default class ListenbrainzScrobbler extends AbstractScrobbleClient { throw e; } } - - getScrobblesForRefresh = async (limit: number) => { - if(this.queuedScrobbles.length === 0) { - return await this.getScrobblesForTimeRange({limit}); - } else { - return await this.getScrobblesForTimeRange({limit, from: this.queuedScrobbles[0].play.data.playDate.unix(), to: dayjs().unix()}); - } - } - + public playToClientPayload(playObj: PlayObject): ListenPayload { return playToListenPayload(playObj); } diff --git a/src/backend/scrobblers/MalojaScrobbler.ts b/src/backend/scrobblers/MalojaScrobbler.ts index d3a76432..6f9b235a 100644 --- a/src/backend/scrobblers/MalojaScrobbler.ts +++ b/src/backend/scrobblers/MalojaScrobbler.ts @@ -79,15 +79,7 @@ export default class MalojaScrobbler extends AbstractScrobbleClient { throw e; } } - - getScrobblesForRefresh = async (limit: number) => { - if(this.queuedScrobbles.length === 0) { - return await this.getScrobblesForTimeRange({limit, cursor: 0}); - } else { - return await this.getScrobblesForTimeRange({limit, cursor: 0, from: this.queuedScrobbles[0].play.data.playDate.unix(), to: dayjs().unix()}); - } - } - + public playToClientPayload(playObj: PlayObject): MalojaScrobbleRequestData { const { apiKey } = this.config.data; diff --git a/src/backend/scrobblers/RockskyScrobbler.ts b/src/backend/scrobblers/RockskyScrobbler.ts index d18e3204..26b2ff82 100644 --- a/src/backend/scrobblers/RockskyScrobbler.ts +++ b/src/backend/scrobblers/RockskyScrobbler.ts @@ -66,8 +66,8 @@ export default class RockskyScrobbler extends AbstractScrobbleClient { } } - getScrobblesForRefresh = async (limit: number) => { - return await this.api.getRecentlyPlayed(limit); + getScrobblesForTimeRange = async (_) => { + return await this.api.getRecentlyPlayed(this.MAX_INITIAL_SCROBBLES_FETCH); } public playToClientPayload(playObj: PlayObject): ListenPayload { diff --git a/src/backend/scrobblers/TealfmScrobbler.ts b/src/backend/scrobblers/TealfmScrobbler.ts index 0c7613e2..72b0d60f 100644 --- a/src/backend/scrobblers/TealfmScrobbler.ts +++ b/src/backend/scrobblers/TealfmScrobbler.ts @@ -91,9 +91,9 @@ export default class TealScrobbler extends AbstractScrobbleClient { } } - getScrobblesForRefresh = async (limit: number) => { + getScrobblesForTimeRange = async (_) => { try { - const {data} = await this.client.getPagelessTimeRangeListens({limit}) + const {data} = await this.client.getPagelessTimeRangeListens({limit: 100}) return data; } catch (e) { throw new Error('Error occurred while trying to fetch records', {cause: e}); diff --git a/src/backend/tests/scrobbler/TestScrobbler.ts b/src/backend/tests/scrobbler/TestScrobbler.ts index 75a66ac1..d7847c57 100644 --- a/src/backend/tests/scrobbler/TestScrobbler.ts +++ b/src/backend/tests/scrobbler/TestScrobbler.ts @@ -6,20 +6,21 @@ import { Notifiers } from "../../notifier/Notifiers.js"; import AbstractScrobbleClient from "../../scrobblers/AbstractScrobbleClient.js"; import { CommonClientConfig, CommonClientOptions, NowPlayingOptions } from "../../common/infrastructure/config/client/index.js"; import clone from "clone"; +import { TimeRangeListensFetcher } from "../../common/infrastructure/Atomic.js"; export class TestScrobbler extends AbstractScrobbleClient { testRecentScrobbles: PlayObject[] = []; + getScrobblesForTimeRange: TimeRangeListensFetcher; constructor(config: CommonClientConfig = {name: 'test'}) { const logger = loggerTest; const notifier = new Notifiers(new EventEmitter(), new EventEmitter(), new EventEmitter(), logger); super('test', 'Test', {name: 'test', ...config}, notifier, new EventEmitter(), logger); this.supportsNowPlaying = false; - } - - protected async getScrobblesForRefresh(limit: number): Promise { - return this.testRecentScrobbles; + this.getScrobblesForTimeRange = async (_) => this.testRecentScrobbles; + this.scrobbleDelay = 10; + this.scrobbleSleep = 20; } doScrobble(playObj: PlayObject) { diff --git a/src/backend/tests/scrobbler/scrobblers.test.ts b/src/backend/tests/scrobbler/scrobblers.test.ts index 709afac3..9796db95 100644 --- a/src/backend/tests/scrobbler/scrobblers.test.ts +++ b/src/backend/tests/scrobbler/scrobblers.test.ts @@ -1,13 +1,13 @@ import chai, { assert, expect } from 'chai'; +import {spy} from 'sinon'; import asPromised from 'chai-as-promised'; import clone from 'clone'; -import { source } from "common-tags"; import dayjs from "dayjs"; import { after, before, describe, it } from 'mocha'; import { http, HttpResponse } from 'msw'; import pEvent from 'p-event'; import { PlayObject } from "../../../core/Atomic.js"; -import { genGroupIdStr, sleep } from "../../utils.js"; +import { genGroupIdStr, sleep, sortByOldestPlayDate } from "../../utils.js"; import mixedDuration from '../plays/mixedDuration.json' with { type: 'json' }; import withDuration from '../plays/withDuration.json' with { type: 'json' }; import { MockNetworkError, withRequestInterception } from "../utils/networking.js"; @@ -15,8 +15,10 @@ import { asPlays, generatePlay, generatePlayPlatformId, generatePlays, generateS import MockDate from 'mockdate'; import { NowPlayingScrobbler, TestAuthScrobbler, TestScrobbler } from "./TestScrobbler.js"; -import { PlayPlatformId } from '../../common/infrastructure/Atomic.js'; +import { PaginatedTimeRangeOptions, PlayPlatformId, REFRESH_STALE_DEFAULT } from '../../common/infrastructure/Atomic.js'; import { defaultLifecycle } from '../../utils/PlayTransformUtils.js'; +import { shuffleArray } from '../../utils/DataUtils.js'; +import { DEFAULT_GROUP_DURATION, groupPlaysToTimeRanges } from '../../utils/ListenFetchUtils.js'; chai.use(asPromised); @@ -39,7 +41,6 @@ const generateTestScrobbler = () => { confidenceBreakdown: true } }; - testScrobbler.lastScrobbleCheck = dayjs().subtract(60, 'seconds'); return testScrobbler; } @@ -117,7 +118,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu it('It is not detected as duplicate when play date is newer than most recent', async function () { - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; const newScrobble = generatePlay({ playDate: normalizedWithMixedDur[normalizedWithMixedDur.length - 1].data.playDate.add(70, 'seconds') @@ -129,7 +130,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu it('It is not detected as duplicate when play date is close to an existing scrobble', async function () { - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; const newScrobble = generatePlay({ playDate: normalizedWithMixedDur[normalizedWithMixedDur.length - 3].data.playDate.add(3, 'seconds') @@ -140,7 +141,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu it('It handles unique detection when no existing scrobble matches above a score of 0', async function () { - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; const uniquePlay = generatePlay({ artists: [ @@ -164,7 +165,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu it('Is not detected as duplicate when artist is same, time is similar, but track is different', async function () { - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; const diffPlay = clone(normalizedWithMixedDur[1]); diffPlay.data.playDate = diffPlay.data.playDate.add(9, 's'); @@ -175,7 +176,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu it('Is not detected as duplicate when track is same, time is similar, but artist is different', async function () { - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; const diffPlay = clone(normalizedWithMixedDur[1]); diffPlay.data.playDate = diffPlay.data.playDate.add(9, 's'); @@ -187,7 +188,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu it('Is not detected as duplicate when play date is different by more than 10 seconds (high granularity source)', async function () { - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; const timeOffPos = clone(normalizedWithMixedDur[normalizedWithMixedDur.length - 1]); timeOffPos.data.playDate = timeOffPos.data.playDate.add(11, 's'); @@ -205,7 +206,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu initialDate: firstPlayDate, defaultMeta: {source: 'subsonic'} }); - testScrobbler.recentScrobbles = recent; + testScrobbler.testRecentScrobbles = recent; const timeOffPos = clone(recent[recent.length - 1]); timeOffPos.data.playDate = timeOffPos.data.playDate.add(61, 's'); @@ -225,7 +226,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu it('A track with continuity to the previous track is not detected as a duplicate', async function () { - testScrobbler.recentScrobbles = normalizedWithDur; + testScrobbler.testRecentScrobbles = normalizedWithDur; const brickPt1 = normalizedWithDur.find(x => x.data.track.includes('Another Brick')); const brickPt2 = clone(brickPt1); @@ -254,7 +255,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu initialDate: firstPlayDate, defaultMeta: {source: 'jellyfin'} }); - testScrobbler.recentScrobbles = recent; + testScrobbler.testRecentScrobbles = recent; const repeatPlay = clone(recent[recent.length - 1]); repeatPlay.data.playDate = repeatPlay.data.playDate.add(repeatPlay.data.duration + 2, 's'); @@ -275,12 +276,12 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu }); it('Is detected as duplicate when an exact match', async function () { - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; assert.isTrue((await testScrobbler.alreadyScrobbled(normalizedWithMixedDur[normalizedWithMixedDur.length - 1]))[0]); }); it('Is detected as duplicate when artist/title differences are whitespace or case', async function () { - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; const ref = normalizedWithMixedDur[3]; const diffPlay = clone(ref); @@ -304,7 +305,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu }); it('Is detected as duplicate when artist/title differences are from unicode normalization', async function () { - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; const ref = normalizedWithMixedDur.find(x => x.data.track === 'Jimbó'); const diffPlay = clone(ref); @@ -315,7 +316,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu it('Is detected as duplicate when play date is off by 10 seconds or less (high granularity source)', async function () { - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; const timeOffPos = clone(normalizedWithMixedDur[normalizedWithMixedDur.length - 1]); timeOffPos.data.playDate = timeOffPos.data.playDate.add(10, 's'); @@ -331,7 +332,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu son.data.playDate = dayjs().subtract(1, 'hour').set('minute', 26).set('second', 20); son.data.duration = 267; son.data.listenedFor = undefined; - testScrobbler.recentScrobbles = normalizedWithMixedDurOlder.concat(son); + testScrobbler.testRecentScrobbles = normalizedWithMixedDurOlder.concat(son); const offSon = clone(son); offSon.data.playDate = dayjs().subtract(1, 'hour').set('minute', 30).set('second', 37); @@ -344,7 +345,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu initialDate: firstPlayDate, defaultMeta: {source: 'subsonic'} }); - testScrobbler.recentScrobbles = recent; + testScrobbler.testRecentScrobbles = recent; const timeOffPos = clone(recent[recent.length - 1]); timeOffPos.data.playDate = timeOffPos.data.playDate.add(59, 's'); @@ -357,7 +358,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu }); it('Is detected as duplicate when title is exact, artist is similar, and time is similar', async function () { - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; const ref = normalizedWithMixedDur[3]; const diffPlay = clone(ref); @@ -402,7 +403,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu } } - testScrobbler.recentScrobbles = normalizedWithMixedDurOlder.concat(ref); + testScrobbler.testRecentScrobbles = normalizedWithMixedDurOlder.concat(ref); assert.isTrue((await testScrobbler.alreadyScrobbled(spotifyPlay))[0]); }); @@ -415,7 +416,7 @@ describe('Detects duplicate and unique scrobbles from client recent history', fu it('Is detected as duplicate when play date is close to the end of an existing scrobble', async function () { - testScrobbler.recentScrobbles = normalizedWithDur; + testScrobbler.testRecentScrobbles = normalizedWithDur; const timeEnd = clone(normalizedWithDur[normalizedWithMixedDur.length - 2]); timeEnd.data.playDate = timeEnd.data.playDate.add(timeEnd.data.duration, 's'); @@ -441,15 +442,9 @@ describe('Detects duplicate and unique scrobbles using actively tracked scrobble beforeEach(function() { testScrobbler = generateTestScrobbler(); - testScrobbler.recentScrobbles = normalizedWithMixedDur; - testScrobbler.lastScrobbleCheck = dayjs().subtract(60, 'seconds'); + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; }); - // before(function () { - // testScrobbler.recentScrobbles = normalizedWithMixedDur; - // testScrobbler.lastScrobbleCheck = dayjs().subtract(60, 'seconds'); - // }); - it('Detects a unique play', async function() { const newScrobble = generatePlay({ playDate: normalizedWithMixedDur[normalizedWithMixedDur.length - 3].data.playDate.add(3, 'seconds') @@ -491,65 +486,86 @@ describe('Detects duplicate and unique scrobbles using actively tracked scrobble describe('Upstream Scrobbles', function() { - it('Stores upstream scrobbles on refresh', async function () { - const scrobbler = generateTestScrobbler(); - scrobbler.testRecentScrobbles = normalizedWithMixedDur; - assert.isEmpty(scrobbler.recentScrobbles); - await scrobbler.refreshScrobbles(); - assert.isNotEmpty(scrobbler.recentScrobbles); + afterEach(function () { + MockDate.reset(); }); - describe('Detects when upstream scrobbles should be refreshed', function() { - - const normalizedClose = normalizePlays(generatePlays(10), {endDate: dayjs().subtract(100, 'seconds')}); - - beforeEach(async function () { - testScrobbler = generateTestScrobbler(); - testScrobbler.testRecentScrobbles = normalizedClose; - await testScrobbler.initialize(); - testScrobbler.lastScrobbleCheck = dayjs().subtract(65, 'seconds'); - testScrobbler.queuedScrobbles = []; - testScrobbler.config.options = {}; - }); - - it('Detects queued scrobble date is newer than last scrobble refresh', async function() { - const newScrobble = generatePlay({ - playDate: dayjs() - }); - - await testScrobbler.queueScrobble(newScrobble, 'test'); - assert.isTrue(testScrobbler.shouldRefreshScrobble()); - }); - - it('Detects queued scrobble date is older than newest scrobble', async function() { - const newScrobble = generatePlay({ - playDate: dayjs().subtract(120, 'seconds') - }); - - await testScrobbler.queueScrobble(newScrobble, 'test'); - assert.isTrue(testScrobbler.shouldRefreshScrobble()); - }); - - it('Forces refresh if refreshStaleAfter is set', async function() { - testScrobbler.config.options = { refreshStaleAfter: 10 }; - - const newScrobble = generatePlay({ - playDate: dayjs().subtract(80, 'seconds') - }); + it('Calls timerange func to get SOT scrobbles when none exists', async function() { + const existingPlays = normalizePlays(generatePlays(3), {initialDate: dayjs().subtract(1, 'hour')}); + const scrobbler = generateTestScrobbler(); + scrobbler.testRecentScrobbles = existingPlays; + await scrobbler.tryInitialize(); + + const sp = spy(scrobbler, 'getScrobblesForTimeRange'); + + const play = generatePlay({playDate: dayjs().subtract(60, 's')}); + await scrobbler.queueScrobble(play, 'test'); + const emptied = pEvent(scrobbler.emitter, 'queueEmptied'); + scrobbler.startScrobbling().then(() => null); + await emptied; + scrobbler.tryStopScrobbling().then(() => null); + expect(sp.called).is.true; + }); - await testScrobbler.queueScrobble(newScrobble, 'test'); - assert.isTrue(testScrobbler.shouldRefreshScrobble()); - }); + it('Uses cached timerange for closely grouped scrobbles', async function() { + const existingPlays = normalizePlays(generatePlays(3), {initialDate: dayjs().subtract(1, 'hour')}); + const scrobbler = generateTestScrobbler(); + scrobbler.testRecentScrobbles = existingPlays; + await scrobbler.tryInitialize(); + + const sp = spy(scrobbler, 'getScrobblesForTimeRange'); + + const play1 = generatePlay({playDate: dayjs().subtract(3, 'm')}); + const play2 = generatePlay({playDate: dayjs().subtract(1, 'm')}); + await scrobbler.queueScrobble([play1, play2], 'test'); + const emptied = pEvent(scrobbler.emitter, 'queueEmptied'); + scrobbler.startScrobbling().then(() => null); + await emptied; + scrobbler.tryStopScrobbling().then(() => null); + expect(sp.calledOnce).is.true; + }); - it('Does not refresh if scrobble is older than last check but newer than newest upstream scrobble', async function() { - testScrobbler.lastScrobbleCheck = dayjs().subtract(40, 'seconds'); - const newScrobble = generatePlay({ - playDate: dayjs().subtract(80, 'seconds') - }); + it('Uses separate timerange calls when scrobbles are not closely grouped', async function() { + const existingPlays = normalizePlays(generatePlays(3), {initialDate: dayjs().subtract(1, 'hour')}); + const scrobbler = generateTestScrobbler(); + scrobbler.testRecentScrobbles = existingPlays; + await scrobbler.tryInitialize(); + + const sp = spy(scrobbler, 'getScrobblesForTimeRange'); + + const play1 = generatePlay({playDate: dayjs().subtract(3, 'm')}); + const play2 = generatePlay({playDate: dayjs().subtract(1, 'm')}); + const play3 = generatePlay({playDate: dayjs().subtract(DEFAULT_GROUP_DURATION.add(4, 'm'))}); + await scrobbler.queueScrobble([play1, play2, play3], 'test'); + const emptied = pEvent(scrobbler.emitter, 'queueEmptied'); + scrobbler.startScrobbling().then(() => null); + await emptied; + scrobbler.tryStopScrobbling().then(() => null); + expect(sp.calledTwice).is.true; + }); - await testScrobbler.queueScrobble(newScrobble, 'test'); - assert.isFalse(testScrobbler.shouldRefreshScrobble()); - }); + it('Gets fresh timerange if TTL of staleAfter has passed', async function() { + const existingPlays = normalizePlays(generatePlays(3), {initialDate: dayjs().subtract(1, 'hour')}); + const scrobbler = generateTestScrobbler(); + scrobbler.testRecentScrobbles = existingPlays; + await scrobbler.tryInitialize(); + + const sp = spy(scrobbler, 'getScrobblesForTimeRange'); + + const play1 = generatePlay({playDate: dayjs().subtract(3, 'm')}); + const play2 = generatePlay({playDate: dayjs().subtract(1, 'm')}); + await scrobbler.queueScrobble([play1], 'test'); + const emptied = pEvent(scrobbler.emitter, 'queueEmptied'); + scrobbler.startScrobbling().then(() => null); + await emptied; + expect(sp.calledOnce).is.true; + + MockDate.set(dayjs().add(REFRESH_STALE_DEFAULT + 1, 's').toDate()); + const emptied2 = pEvent(scrobbler.emitter, 'queueEmptied'); + await scrobbler.queueScrobble([play2], 'test'); + await emptied2; + scrobbler.tryStopScrobbling().then(() => null); + expect(sp.calledTwice).is.true; }); }); @@ -561,10 +577,10 @@ describe('Scrobble client uses transform plays correctly', function() { beforeEach(async function() { testScrobbler = generateTestScrobbler(); await testScrobbler.initialize(); - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; testScrobbler.scrobbleSleep = 500; testScrobbler.scrobbleDelay = 0; - testScrobbler.lastScrobbleCheck = dayjs().subtract(60, 'seconds'); + //testScrobbler.lastScrobbleCheck = dayjs().subtract(60, 'seconds'); testScrobbler.config.options = {}; //testScrobbler.initScrobbleMonitoring().catch(console.error); }); @@ -629,7 +645,7 @@ describe('Scrobble client uses transform plays correctly', function() { track: 'my hugely cool and very different track title' }); - testScrobbler.recentScrobbles = normalizePlays([newScrobble, ...withDurPlays], {initialDate: firstPlayDate}); + testScrobbler.testRecentScrobbles = normalizePlays([newScrobble, ...withDurPlays], {initialDate: firstPlayDate}); testScrobbler.buildTransformRules(); expect((await testScrobbler.alreadyScrobbled(newScrobble))[0]).is.false; @@ -651,7 +667,7 @@ describe('Scrobble client uses transform plays correctly', function() { track: 'my hugely cool and very different track title' }); - testScrobbler.recentScrobbles = normalizePlays([newScrobble, ...withDurPlays], {initialDate: firstPlayDate}); + testScrobbler.testRecentScrobbles = normalizePlays([newScrobble, ...withDurPlays], {initialDate: firstPlayDate}); testScrobbler.buildTransformRules(); expect((await testScrobbler.alreadyScrobbled(newScrobble))[0]).is.false; @@ -669,11 +685,10 @@ describe('Manages scrobble queue', function() { beforeEach(async function() { testScrobbler = generateTestScrobbler(); await testScrobbler.initialize(); - testScrobbler.recentScrobbles = normalizedWithMixedDur; + testScrobbler.testRecentScrobbles = normalizedWithMixedDur; testScrobbler.testRecentScrobbles = normalizedWithMixedDur; testScrobbler.scrobbleSleep = 500; testScrobbler.scrobbleDelay = 0; - testScrobbler.lastScrobbleCheck = dayjs().subtract(60, 'seconds'); testScrobbler.initScrobbleMonitoring().catch(console.error); }); @@ -993,3 +1008,49 @@ describe('Now Playing', function() { }); }); + +describe('Scrobble Temporal Grouping', function () { + + it('Groups into separate groups when not within duration', function() { + const plays1 = normalizePlays(generatePlays(3), {initialDate: dayjs().subtract(1, 'hour')}); + plays1.sort(sortByOldestPlayDate); + const oldest1 = plays1[0].data.playDate.unix(); + const newest1 = plays1[plays1.length - 1].data.playDate.unix(); + + const plays2 = normalizePlays(generatePlays(3), {initialDate: dayjs().subtract(2, 'hour')}); + plays2.sort(sortByOldestPlayDate); + const oldest2 = plays2[0].data.playDate.unix(); + const newest2 = plays2[plays2.length - 1].data.playDate.unix(); + + const plays = [...plays1, ...plays2]; + shuffleArray(plays); + + const ranges = groupPlaysToTimeRanges(plays, []); + expect(ranges.length).eq(2); + expect(ranges.some(x => x.from < oldest1 && x.to > newest1)).is.true; + expect(ranges.some(x => x.from < oldest2 && x.to > newest2)).is.true; + }); + + it('Groups into existing time range', function() { + const plays1 = normalizePlays(generatePlays(3), {initialDate: dayjs().subtract(1, 'hour')}); + plays1.sort(sortByOldestPlayDate); + const oldest1 = plays1[0].data.playDate; + const newest1 = plays1[plays1.length - 1].data.playDate; + + const plays2 = normalizePlays(generatePlays(3), {initialDate: dayjs().subtract(2, 'hour')}); + plays2.sort(sortByOldestPlayDate); + const oldest2 = plays2[0].data.playDate; + const newest2 = plays2[plays2.length - 1].data.playDate; + + const plays = [...plays1, ...plays2]; + shuffleArray(plays); + + const existing: PaginatedTimeRangeOptions = {from: oldest1.subtract(10, 's').unix(), to: newest1.add(10, 's').unix()}; + + const ranges = groupPlaysToTimeRanges(plays, [existing]); + expect(ranges.length).eq(2); + expect(ranges.some(x => x.from < oldest1.unix() && x.to > newest1.unix())).is.true; + expect(ranges.some(x => x.from < oldest2.unix() && x.to > newest2.unix())).is.true; + expect(ranges.some(x => x.to === existing.to && x.from && existing.from)); + }); +}) \ No newline at end of file diff --git a/src/backend/utils/DataUtils.ts b/src/backend/utils/DataUtils.ts index 487fc382..fbc0faef 100644 --- a/src/backend/utils/DataUtils.ts +++ b/src/backend/utils/DataUtils.ts @@ -163,4 +163,15 @@ export const objectIsEmpty = (obj: object, valueIsEmpty?: undefined | ((val: any } return Object.values(obj).every(x => valueIsEmpty(x)); +} + +/** Shuffle array in-place + * + * https://stackoverflow.com/a/12646864/1469797 + */ +export const shuffleArray = (array: any[]): void => { + for (let i = array.length - 1; i > 0; i--) { + const j = Math.floor(Math.random() * (i + 1)); + [array[i], array[j]] = [array[j], array[i]]; + } } \ No newline at end of file diff --git a/src/backend/utils/ListenFetchUtils.ts b/src/backend/utils/ListenFetchUtils.ts index e11b1a39..89af346e 100644 --- a/src/backend/utils/ListenFetchUtils.ts +++ b/src/backend/utils/ListenFetchUtils.ts @@ -1,10 +1,12 @@ -import { childLogger, Logger } from "@foxxmd/logging"; -import dayjs from "dayjs"; -import { PlayObject } from "../../core/Atomic.js"; -import { hasPagelessTimeRangeListens, hasPaginatedTimeRangeListens, PagelessListensTimeRangeOptions, PagelessTimeRangeListens, PaginatedListensTimeRangeOptions, PaginatedTimeRangeCommonOptions, PaginatedTimeRangeListens, PaginatedTimeRangeSource, TimeRangeListensFetcher } from "../common/infrastructure/Atomic.js"; +import { childLogger, Logger, loggerTest } from "@foxxmd/logging"; +import dayjs, { Dayjs } from "dayjs"; +import { Duration } from "dayjs/plugin/duration.js"; +import { PlayObject, UnixTimestamp } from "../../core/Atomic.js"; +import { hasPagelessTimeRangeListens, hasPaginatedTimeRangeListens, PagelessListensTimeRangeOptions, PagelessTimeRangeListens, PaginatedListensTimeRangeOptions, PaginatedTimeRangeCommonOptions, PaginatedTimeRangeListens, PaginatedTimeRangeOptions, PaginatedTimeRangeSource, REFRESH_STALE_DEFAULT, TimeRangeListensFetcher } from "../common/infrastructure/Atomic.js"; import { MaybeLogger } from "../common/logging.js"; import { sortByNewestPlayDate, sortByOldestPlayDate } from "../utils.js"; import { todayAwareFormat } from "./TimeUtils.js"; +import { playDateWithinDurationOfAny } from "./PlayComparisonUtils.js"; export interface TimeRangeFetchOptions { logger?: MaybeLogger | Logger @@ -171,4 +173,65 @@ export const createGetScrobblesForTimeRangeFunc = { + const { + groupDuration = DEFAULT_GROUP_DURATION, + newPadding = DEFAULT_NEW_PADDING, + staleNowBuffer = REFRESH_STALE_DEFAULT, + logger = loggerTest + } = opts; + const newRanges: PaginatedTimeRangeOptions[] = []; + + const temporallyClosePlaySets: PlayObject[][] = []; + + const sorted = [...plays]; + sorted.sort(sortByOldestPlayDate); + + for(const p of sorted) { + const closePlaySetIndex = temporallyClosePlaySets.findIndex(x => playDateWithinDurationOfAny(p, x, groupDuration)); + if(closePlaySetIndex === -1) { + temporallyClosePlaySets.push([p]); + } else { + temporallyClosePlaySets[closePlaySetIndex].push(p); + } + } + + temporallyClosePlaySets.forEach((x) => x.sort(sortByOldestPlayDate)); + for(const tc of temporallyClosePlaySets) { + let oldest: Dayjs, + newest: Dayjs; + if(tc.length === 1) { + oldest = tc[0].data.playDate; + newest = oldest; + //newest = tc[0].data.playDate.add(1, 'hour').unix(); + } else { + oldest = tc[0].data.playDate; + newest = tc[tc.length - 1].data.playDate; + } + + let bufferedNewest = newest; + if(dayjs().diff(bufferedNewest, 's') < staleNowBuffer) { + bufferedNewest = bufferedNewest.subtract(staleNowBuffer, 's'); + } + const existingWithin = existingRanges.find(x => x.from <= oldest.unix() && x.to >= bufferedNewest.unix()); + if(!existingWithin) { + newRanges.push({from: oldest.subtract(newPadding).unix(), to: Math.min(newest.add(newPadding).unix(), dayjs().unix())}); + } else { + newRanges.push(existingWithin); + } + } + + return newRanges; } \ No newline at end of file diff --git a/src/backend/utils/PlayComparisonUtils.ts b/src/backend/utils/PlayComparisonUtils.ts index ec27a3de..e8502670 100644 --- a/src/backend/utils/PlayComparisonUtils.ts +++ b/src/backend/utils/PlayComparisonUtils.ts @@ -6,6 +6,7 @@ import { comparePlayTemporally, hasAcceptableTemporalAccuracy, TemporalPlayCompa import { compareNormalizedStrings, compareScrobbleArtists, compareScrobbleTracks, compareTracks, normalizeStr, TrackSamenessResults } from "./StringUtils.js"; import { ARTIST_WEIGHT, TITLE_WEIGHT } from "../common/infrastructure/Atomic.js"; import { StringSamenessResult } from "@foxxmd/string-sameness"; +import { Duration } from "dayjs/plugin/duration.js"; export const metaInvariantTransform = (play: PlayObject): PlayObjectLifecycleless => { @@ -377,4 +378,8 @@ export const scorePlaySameness = (ref: PlayObject, candidate: PlayObject, option const albumScore = albumHigh * (albumWeight + albumBonus); return trackScore + artistScore + albumScore; +} + +export const playDateWithinDurationOfAny = (play: PlayObject, plays: PlayObject[], dur: Duration): PlayObject | undefined => { + return plays.find(x => Math.abs(x.data.playDate.diff(play.data.playDate, 's')) <= dur.asSeconds()); } \ No newline at end of file -- 2.51.2