From b3217f2d9c4561afc06034023671a77bf521409b Mon Sep 17 00:00:00 2001 From: FoxxMD Date: Wed, 11 Feb 2026 21:39:21 +0000 Subject: [PATCH] mostly implemented time range fetch utility --- src/backend/common/infrastructure/Atomic.ts | 10 +- .../config/client/listenbrainz.ts | 8 + .../common/vendor/ListenbrainzApiClient.ts | 14 +- src/backend/scrobblers/KoitoScrobbler.ts | 14 +- src/backend/scrobblers/LastfmScrobbler.ts | 12 +- .../scrobblers/ListenbrainzScrobbler.ts | 16 +- src/backend/scrobblers/MalojaScrobbler.ts | 13 +- src/backend/sources/LastfmSource.ts | 7 +- src/backend/utils.ts | 3 +- src/backend/utils/ListenFetchUtils.ts | 162 ++++++++++++++++++ 10 files changed, 229 insertions(+), 30 deletions(-) create mode 100644 src/backend/utils/ListenFetchUtils.ts diff --git a/src/backend/common/infrastructure/Atomic.ts b/src/backend/common/infrastructure/Atomic.ts index 338be622..93f4b726 100644 --- a/src/backend/common/infrastructure/Atomic.ts +++ b/src/backend/common/infrastructure/Atomic.ts @@ -443,11 +443,13 @@ export interface PaginatedTimeRangeOptions { to: UnixTimestamp } +export type PaginatedTimeRangeCommonOptions = Partial & PaginatedLimit; + export interface PaginatedListensOptions extends PaginatedLimit { page: number } -export interface PagelessListensTimeRangeOptions extends Partial, PaginatedLimit { +export interface PagelessListensTimeRangeOptions extends PaginatedTimeRangeCommonOptions { } export interface PaginatedListensTimeRangeOptions extends Partial, PaginatedListensOptions { @@ -456,6 +458,7 @@ export interface PaginatedListensTimeRangeOptions extends Partial Promise \ No newline at end of file diff --git a/src/backend/common/infrastructure/config/client/listenbrainz.ts b/src/backend/common/infrastructure/config/client/listenbrainz.ts index a852db39..0708c6a1 100644 --- a/src/backend/common/infrastructure/config/client/listenbrainz.ts +++ b/src/backend/common/infrastructure/config/client/listenbrainz.ts @@ -38,3 +38,11 @@ export interface ListenBrainzClientConfig extends CommonClientConfig { export interface ListenBrainzClientAIOConfig extends ListenBrainzClientConfig { type: 'listenbrainz' } + + +/** https://github.com/metabrainz/listenbrainz-server/pull/2572 + * https://github.com/metabrainz/listenbrainz-server/blob/master/listenbrainz/webserver/views/api_tools.py#L48 + */ +export const MAX_ITEMS_PER_GET_LZ = 1000; +export const DEFAULT_ITEMS_PER_GET_LZ = 25; +export const DEFAULT_MS_ITEMS_PER_GET_LZ = 100; \ No newline at end of file diff --git a/src/backend/common/vendor/ListenbrainzApiClient.ts b/src/backend/common/vendor/ListenbrainzApiClient.ts index 504a21be..3a8a6698 100644 --- a/src/backend/common/vendor/ListenbrainzApiClient.ts +++ b/src/backend/common/vendor/ListenbrainzApiClient.ts @@ -15,7 +15,7 @@ import { import { getScrobbleTsSOCDate } from "../../utils/TimeUtils.js"; import { UpstreamError } from "../errors/UpstreamError.js"; import { AbstractApiOptions, DEFAULT_RETRY_MULTIPLIER, DELIMITERS, FormatPlayObjectOptions, PagelessListensTimeRangeOptions, PagelessTimeRangeListens, PagelessTimeRangeListensResult } from "../infrastructure/Atomic.js"; -import { ListenBrainzClientData } from "../infrastructure/config/client/listenbrainz.js"; +import { DEFAULT_ITEMS_PER_GET_LZ, ListenBrainzClientData, MAX_ITEMS_PER_GET_LZ } from "../infrastructure/config/client/listenbrainz.js"; import AbstractApiClient from "./AbstractApiClient.js"; import { getBaseFromUrl, isPortReachableConnect, joinedUrl, normalizeWebAddress } from '../../utils/NetworkUtils.js'; import { isEmptyArrayOrUndefined, removeUndefinedKeys, unique } from '../../utils.js'; @@ -43,19 +43,13 @@ export interface UserListensOptions { max_ts?: UnixTimestamp /** unix epoch timestamp, listens with listened_at greater than (but not including) this value will be returned. */ min_ts?: UnixTimestamp - /** number of listens to return. Max is `MAX_ITEMS_PER_GET` + /** number of listens to return. Max is `MAX_ITEMS_PER_GET_LZ` * * @default 25 */ count?: number } -/** https://github.com/metabrainz/listenbrainz-server/pull/2572 - * https://github.com/metabrainz/listenbrainz-server/blob/master/listenbrainz/webserver/views/api_tools.py#L48 - */ -const MAX_ITEMS_PER_GET = 1000; -const DEFAULT_ITEMS_PER_GET = 25; - export class ListenbrainzApiClient extends AbstractApiClient implements PagelessTimeRangeListens { declare config: ListenBrainzClientData; @@ -295,7 +289,7 @@ export class ListenbrainzApiClient extends AbstractApiClient implements Pageless response: 15000, deadline: 30000 }) - .query({...options, count: Math.min(count, MAX_ITEMS_PER_GET)})); + .query({...options, count: Math.min(count, DEFAULT_ITEMS_PER_GET_LZ)})); const {body: {payload}} = resp as any; return payload as ListensResponse; @@ -308,7 +302,7 @@ export class ListenbrainzApiClient extends AbstractApiClient implements Pageless const { limit = 100, to, from, user } = options; const lzListensOptions: UserListensOptions = { - count: Math.min(limit, MAX_ITEMS_PER_GET), + count: Math.min(limit, MAX_ITEMS_PER_GET_LZ), }; try { diff --git a/src/backend/scrobblers/KoitoScrobbler.ts b/src/backend/scrobblers/KoitoScrobbler.ts index 82d1d9de..639bf3d3 100644 --- a/src/backend/scrobblers/KoitoScrobbler.ts +++ b/src/backend/scrobblers/KoitoScrobbler.ts @@ -4,7 +4,7 @@ import { PlayObject, SourcePlayerObj } 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 { FormatPlayObjectOptions } from "../common/infrastructure/Atomic.js"; +import { FormatPlayObjectOptions, TimeRangeListensFetcher } from "../common/infrastructure/Atomic.js"; import { playToListenPayload } from "../common/vendor/ListenbrainzApiClient.js"; import { Notifiers } from "../notifier/Notifiers.js"; @@ -12,13 +12,15 @@ import AbstractScrobbleClient, { shouldUpdatePlayingNowPlatformWhenPlayingOnly } import { isDebugMode } from "../utils.js"; import { KoitoClientConfig } from "../common/infrastructure/config/client/koito.js"; import { KoitoApiClient, listenObjectResponseToPlay } from "../common/vendor/koito/KoitoApiClient.js"; +import { createGetScrobblesForTimeRangeFunc } from "../utils/ListenFetchUtils.js"; +import dayjs from "dayjs"; export default class KoitoScrobbler extends AbstractScrobbleClient { api: KoitoApiClient; requiresAuth = true; requiresAuthInteraction = false; - + getScrobblesForTimeRange: TimeRangeListensFetcher declare config: KoitoClientConfig; constructor(name: any, config: KoitoClientConfig, options = {}, notifier: Notifiers, emitter: EventEmitter, logger: Logger) { @@ -28,6 +30,7 @@ export default class KoitoScrobbler extends AbstractScrobbleClient { // 1000 is way too high. maxing at 100 this.MAX_INITIAL_SCROBBLES_FETCH = 100; this.supportsNowPlaying = true; + this.getScrobblesForTimeRange = createGetScrobblesForTimeRangeFunc(this.api, this.api.logger); } formatPlayObj = (obj: any, options: FormatPlayObjectOptions = {}) => listenObjectResponseToPlay(obj, options); @@ -67,8 +70,11 @@ export default class KoitoScrobbler extends AbstractScrobbleClient { } getScrobblesForRefresh = async (limit: number) => { - const resp = await this.api.getPaginatedTimeRangeListens({limit, page: 0 }); - return resp.data; + if(this.queuedScrobbles.length === 0) { + return await this.getScrobblesForTimeRange({limit, page: 0}); + } else { + return await this.getScrobblesForTimeRange({limit, page: 0, from: this.queuedScrobbles[0].play.data.playDate.unix(), to: dayjs().unix()}); + } } doScrobble = async (playObj: PlayObject) => { diff --git a/src/backend/scrobblers/LastfmScrobbler.ts b/src/backend/scrobblers/LastfmScrobbler.ts index 1caabd92..9e4e54a8 100644 --- a/src/backend/scrobblers/LastfmScrobbler.ts +++ b/src/backend/scrobblers/LastfmScrobbler.ts @@ -5,12 +5,13 @@ import { PlayObject, SourcePlayerObj } 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 { FormatPlayObjectOptions, InternalConfigOptional } from "../common/infrastructure/Atomic.js"; +import { FormatPlayObjectOptions, InternalConfigOptional, TimeRangeListensFetcher } from "../common/infrastructure/Atomic.js"; import { LastfmClientConfig } from "../common/infrastructure/config/client/lastfm.js"; import LastfmApiClient, { LastFMIgnoredScrobble, playToClientPayload, formatPlayObj, LASTFM_HOST, LASTFM_PATH } from "../common/vendor/LastfmApiClient.js"; import { Notifiers } from "../notifier/Notifiers.js"; import AbstractScrobbleClient, { nowPlayingUpdateByPlayDuration, shouldUpdatePlayingNowPlatformWhenPlayingOnly } from "./AbstractScrobbleClient.js"; import { findCauseByReference } from "../utils/ErrorUtils.js"; +import { createGetScrobblesForTimeRangeFunc } from "../utils/ListenFetchUtils.js"; export default class LastfmScrobbler extends AbstractScrobbleClient { @@ -18,6 +19,7 @@ export default class LastfmScrobbler extends AbstractScrobbleClient { requiresAuth = true; requiresAuthInteraction = true; upstreamType: string = 'Last.fm'; + getScrobblesForTimeRange: TimeRangeListensFetcher declare config: LastfmClientConfig; @@ -29,6 +31,7 @@ export default class LastfmScrobbler extends AbstractScrobbleClient { this.supportsNowPlaying = true; // last.fm shows Now Playing for the same time as the duration of the track being submitted this.nowPlayingMaxThreshold = nowPlayingUpdateByPlayDuration; + this.getScrobblesForTimeRange = createGetScrobblesForTimeRangeFunc(this.api, this.api.logger); } formatPlayObj = (obj: any, options: FormatPlayObjectOptions = {}) => formatPlayObj(obj, options); @@ -59,8 +62,11 @@ export default class LastfmScrobbler extends AbstractScrobbleClient { } getScrobblesForRefresh = async (limit: number) => { - const {data: plays} = await this.api.getPaginatedTimeRangeListens({limit, page: 1}); - return plays; + if(this.queuedScrobbles.length === 0) { + return await this.getScrobblesForTimeRange({limit, page: 1}); + } else { + return await this.getScrobblesForTimeRange({limit, page: 1, from: this.queuedScrobbles[0].play.data.playDate.unix(), to: dayjs().unix()}); + } } // getScrobblesForTimeRange = async (fromDate?: Dayjs, toDate?: Dayjs, limit: number = 1000): Promise => { diff --git a/src/backend/scrobblers/ListenbrainzScrobbler.ts b/src/backend/scrobblers/ListenbrainzScrobbler.ts index 3d42d8c2..166ca3d2 100644 --- a/src/backend/scrobblers/ListenbrainzScrobbler.ts +++ b/src/backend/scrobblers/ListenbrainzScrobbler.ts @@ -5,14 +5,15 @@ import { PlayObject, SourcePlayerObj } from "../../core/Atomic.js"; import { buildTrackString, capitalize } from "../../core/StringUtils.js"; import { isNodeNetworkException } from "../common/errors/NodeErrors.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 { FormatPlayObjectOptions, TimeRangeListensFetcher } from "../common/infrastructure/Atomic.js"; +import { DEFAULT_MS_ITEMS_PER_GET_LZ, ListenBrainzClientConfig } from "../common/infrastructure/config/client/listenbrainz.js"; import { ListenbrainzApiClient, playToListenPayload, playToSubmitPayload } from "../common/vendor/ListenbrainzApiClient.js"; import { ListenPayload } from '../common/vendor/listenbrainz/interfaces.js'; import { Notifiers } from "../notifier/Notifiers.js"; import AbstractScrobbleClient, { nowPlayingUpdateByPlayDuration, shouldUpdatePlayingNowPlatformWhenPlayingOnly } from "./AbstractScrobbleClient.js"; import { isDebugMode } from "../utils.js"; +import { createGetScrobblesForTimeRangeFunc } from "../utils/ListenFetchUtils.js"; export default class ListenbrainzScrobbler extends AbstractScrobbleClient { @@ -20,7 +21,7 @@ export default class ListenbrainzScrobbler extends AbstractScrobbleClient { requiresAuth = true; requiresAuthInteraction = false; isKoito: boolean = false; - + getScrobblesForTimeRange: TimeRangeListensFetcher declare config: ListenBrainzClientConfig; constructor(name: any, config: ListenBrainzClientConfig, options = {}, notifier: Notifiers, emitter: EventEmitter, logger: Logger) { @@ -28,10 +29,11 @@ export default class ListenbrainzScrobbler extends AbstractScrobbleClient { this.api = new ListenbrainzApiClient(name, config.data, {logger: this.logger}); // https://listenbrainz.readthedocs.io/en/latest/users/api/core.html#get--1-user-(user_name)-listens // 1000 is way too high. maxing at 100 - this.MAX_INITIAL_SCROBBLES_FETCH = 100; + this.MAX_INITIAL_SCROBBLES_FETCH = DEFAULT_MS_ITEMS_PER_GET_LZ; this.supportsNowPlaying = true; // listenbrainz shows Now Playing for the same time as the duration of the track being submitted this.nowPlayingMaxThreshold = nowPlayingUpdateByPlayDuration; + this.getScrobblesForTimeRange = createGetScrobblesForTimeRangeFunc(this.api, this.api.logger); } formatPlayObj = (obj: any, options: FormatPlayObjectOptions = {}) => ListenbrainzApiClient.formatPlayObj(obj, options); @@ -67,7 +69,11 @@ export default class ListenbrainzScrobbler extends AbstractScrobbleClient { } getScrobblesForRefresh = async (limit: number) => { - return await this.api.getRecentlyPlayed(limit); + 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()}); + } } // getScrobblesForTimeRange = async (fromDate?: Dayjs, toDate?: Dayjs, limit: number = 1000): Promise => { diff --git a/src/backend/scrobblers/MalojaScrobbler.ts b/src/backend/scrobblers/MalojaScrobbler.ts index 49d4dced..532d86e9 100644 --- a/src/backend/scrobblers/MalojaScrobbler.ts +++ b/src/backend/scrobblers/MalojaScrobbler.ts @@ -4,7 +4,7 @@ import normalizeUrl from "normalize-url"; import { PlayObject } from "../../core/Atomic.js"; import { buildTrackString, capitalize } from "../../core/StringUtils.js"; import { isNodeNetworkException } from "../common/errors/NodeErrors.js"; -import { FormatPlayObjectOptions } from "../common/infrastructure/Atomic.js"; +import { FormatPlayObjectOptions, TimeRangeListensFetcher } from "../common/infrastructure/Atomic.js"; import { MalojaClientConfig } from "../common/infrastructure/config/client/maloja.js"; import { MalojaScrobbleRequestData, @@ -13,6 +13,8 @@ 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"; +import { createGetScrobblesForTimeRangeFunc } from "../utils/ListenFetchUtils.js"; +import dayjs from "dayjs"; const feat = ["ft.", "ft", "feat.", "feat", "featuring", "Ft.", "Ft", "Feat.", "Feat", "Featuring"]; @@ -23,6 +25,7 @@ export default class MalojaScrobbler extends AbstractScrobbleClient { webUrl: string; api: MalojaApiClient; + getScrobblesForTimeRange: TimeRangeListensFetcher declare config: MalojaClientConfig @@ -30,6 +33,7 @@ export default class MalojaScrobbler extends AbstractScrobbleClient { super('maloja', name, config, notifier, emitter, logger); this.api = new MalojaApiClient(name, this.config.data, { logger: childLogger(this.logger, 'API') }); this.MAX_INITIAL_SCROBBLES_FETCH = 100; + this.getScrobblesForTimeRange = createGetScrobblesForTimeRangeFunc(this.api, this.api.logger); } formatPlayObj = (obj: any, options: FormatPlayObjectOptions = {}) => formatMalojaScrobbleToPlay(obj, { url: this.webUrl }); @@ -77,8 +81,11 @@ export default class MalojaScrobbler extends AbstractScrobbleClient { } getScrobblesForRefresh = async (limit: number) => { - const resp = await this.api.getPaginatedTimeRangeListens({limit, page: 0}); - return resp.data; + if(this.queuedScrobbles.length === 0) { + return await this.getScrobblesForTimeRange({limit, page: 0}); + } else { + return await this.getScrobblesForTimeRange({limit, page: 0, from: this.queuedScrobbles[0].play.data.playDate.unix(), to: dayjs().unix()}); + } } public playToClientPayload(playObj: PlayObject): MalojaScrobbleRequestData { diff --git a/src/backend/sources/LastfmSource.ts b/src/backend/sources/LastfmSource.ts index c8ae792f..93b47077 100644 --- a/src/backend/sources/LastfmSource.ts +++ b/src/backend/sources/LastfmSource.ts @@ -3,7 +3,7 @@ import EventEmitter from "events"; import request from "superagent"; import { PlayObject, SOURCE_SOT } from "../../core/Atomic.js"; import { isNodeNetworkException } from "../common/errors/NodeErrors.js"; -import { FormatPlayObjectOptions, InternalConfig, PaginatedListensTimeRangeOptions, PaginatedTimeRangeListens, PlayPlatformId, SourceType } from "../common/infrastructure/Atomic.js"; +import { FormatPlayObjectOptions, InternalConfig, PaginatedListensTimeRangeOptions, PaginatedTimeRangeListens, PlayPlatformId, SourceType, TimeRangeListensFetcher } from "../common/infrastructure/Atomic.js"; import { LastfmSourceConfig } from "../common/infrastructure/config/source/lastfm.js"; import LastfmApiClient, { formatPlayObj } from "../common/vendor/LastfmApiClient.js"; import { sortByOldestPlayDate } from "../utils.js"; @@ -12,6 +12,7 @@ import MemorySource from "./MemorySource.js"; import { Logger } from "@foxxmd/logging"; import { PlayerStateOptions } from "./PlayerState/AbstractPlayerState.js"; import { NowPlayingPlayerState } from "./PlayerState/NowPlayingPlayerState.js"; +import { createGetScrobblesForTimeRangeFunc } from "../utils/ListenFetchUtils.js"; export default class LastfmSource extends MemorySource implements PaginatedTimeRangeListens { @@ -19,6 +20,7 @@ export default class LastfmSource extends MemorySource implements PaginatedTimeR requiresAuth = true; requiresAuthInteraction = true; upstreamType: string = 'Last.fm'; + getScrobblesForTimeRange: TimeRangeListensFetcher declare config: LastfmSourceConfig; @@ -39,7 +41,8 @@ export default class LastfmSource extends MemorySource implements PaginatedTimeR this.playerSourceOfTruth = SOURCE_SOT.HISTORY; // https://www.last.fm/api/show/user.getRecentTracks this.SCROBBLE_BACKLOG_COUNT = 200; - this.logger.info(`Note: The player for this source is an analogue for the 'Now Playing' status exposed by ${this.type} which is NOT used for scrobbling. Instead, the 'recently played' or 'history' information provided by this source is used for scrobbles.`) + this.logger.info(`Note: The player for this source is an analogue for the 'Now Playing' status exposed by ${this.type} which is NOT used for scrobbling. Instead, the 'recently played' or 'history' information provided by this source is used for scrobbles.`); + this.getScrobblesForTimeRange = createGetScrobblesForTimeRangeFunc(this.api, this.api.logger); } static formatPlayObj(obj: any, options: FormatPlayObjectOptions = {}): PlayObject { diff --git a/src/backend/utils.ts b/src/backend/utils.ts index 7f99402b..16cfc70a 100644 --- a/src/backend/utils.ts +++ b/src/backend/utils.ts @@ -63,7 +63,7 @@ export function sleep(ms: any) { return new Promise(resolve => setTimeout(resolve, ms)); } -// sorts playObj formatted objects by playDate in ascending (oldest first) order +/** sorts playObj formatted objects by playDate in ascending (oldest first) order */ export const sortByOldestPlayDate = (a: PlayObject, b: PlayObject) => { const { data: { @@ -87,6 +87,7 @@ export const sortByOldestPlayDate = (a: PlayObject, b: PlayObject) => { return aPlayDate.isAfter(bPlayDate) ? 1 : -1 }; +/** sorts playObj formatted objects by playDate in descending (newest first) order */ export const sortByNewestPlayDate = (a: PlayObject, b: PlayObject) => { const { data: { diff --git a/src/backend/utils/ListenFetchUtils.ts b/src/backend/utils/ListenFetchUtils.ts new file mode 100644 index 00000000..1c8a4dd2 --- /dev/null +++ b/src/backend/utils/ListenFetchUtils.ts @@ -0,0 +1,162 @@ +import { childLogger } 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 { MaybeLogger } from "../common/logging.js"; +import { sortByNewestPlayDate, sortByOldestPlayDate } from "../utils.js"; +import { todayAwareFormat } from "./TimeUtils.js"; + +export const createGetScrobblesForTimeRangeFunc = (fetcher: T, pLogger: MaybeLogger = new MaybeLogger()): TimeRangeListensFetcher => { + let requestCount: number; + const logger = pLogger instanceof MaybeLogger ? pLogger : childLogger(pLogger, ['Pagination']); + const reqLabel = () => `Request ${requestCount}`; + const reqLogger = logger instanceof MaybeLogger ? pLogger : childLogger(logger, [reqLabel]); + + let plays: PlayObject[] = []; + + if (hasPagelessTimeRangeListens(fetcher)) { + return async (opts: PaginatedTimeRangeCommonOptions): Promise => { + requestCount = 0; + let more = true; + let currOpts = { ...opts }; + let initial = true; + while (more) { + requestCount++; + const reqOptsHint: string[] = [ + `Between ${todayAwareFormat(dayjs(currOpts.from))} and ${todayAwareFormat(dayjs(currOpts.to))}` + ]; + if(currOpts.to !== undefined && currOpts.from !== undefined) { + reqOptsHint.push(`Between ${todayAwareFormat(dayjs(currOpts.from))} and ${todayAwareFormat(dayjs(currOpts.to))}`); + } else if(currOpts.to) { + reqOptsHint.push(`Until ${todayAwareFormat(dayjs(currOpts.to))}`); + } else if(currOpts.to) { + reqOptsHint.push(`From ${todayAwareFormat(dayjs(currOpts.from))}`); + } + + if(currOpts.limit !== undefined) { + reqOptsHint.push(`Limit ${currOpts.limit}`); + } + reqLogger.debug(`Fetching => ${reqOptsHint.join(' | ')}`); + const results = await fetcher.getPagelessTimeRangeListens(currOpts); + if(initial) { + initial = false; + const initialFetchLog = []; + if(results.meta.total !== undefined) { + initialFetchLog.push(`API reported ${results.meta.total} total results`); + } + if(results.meta.limit !== undefined && results.meta.limit !== currOpts.limit) { + initialFetchLog.push(`API reported new limit ${results.meta.limit}`); + currOpts.limit = results.meta.limit; + } + if(initialFetchLog.length > 0) { + logger.debug(initialFetchLog.join(' | ')); + } + } + reqLogger.debug(`${results.data.length} results returned${results.data.length === 0 ? ', ending fetch' : ''}`); + plays = plays.concat(results.data); + if (!results.meta.more) { + logger.debug('API indicated no more results, ending fetch'); + more = false; + } + if (results.data.length === 0) { + more = false; + } + // failsafe? + if(more && results.meta.limit !== undefined && results.data.length < results.meta.limit) { + reqLogger.debug(`Number of returned results was less than reported/defined limit (${results.meta.limit}), ending fetch`); + more = false; + } + + if(more && opts.to === undefined && opts.from === undefined) { + // only wanted one fetch + logger.debug('No to/from defined, ending fetch'); + more = false; + } + + if(more) { + if (results.meta.order === undefined || results.meta.order === 'asc') { + // if meta.order is ascending then assumption the response returns *oldest first* list + // so that the newest play from the response should be used as the new `from` + const nextFrom = [...results.data].sort(sortByNewestPlayDate)[0].data.playDate.unix() + 1; + currOpts.from = nextFrom; + } else { + // otherwise, oldest found play should be the new `to` + const nextTo = [...results.data].sort(sortByOldestPlayDate)[0].data.playDate.unix() - 1; + currOpts.to = nextTo; + } + } + } + return plays; + } + } else if (hasPaginatedTimeRangeListens(fetcher)) { + return async (opts: PaginatedTimeRangeCommonOptions | PaginatedListensTimeRangeOptions): Promise => { + requestCount = 0; + let more = true; + let currOpts: PaginatedListensTimeRangeOptions = { page: 1, ...opts }; + let initial = true; + let timeRangeHint: string; + if(currOpts.to !== undefined && currOpts.from !== undefined) { + timeRangeHint = `Between ${todayAwareFormat(dayjs.unix(currOpts.from))} and ${todayAwareFormat(dayjs.unix(currOpts.to))}`; + } else if(currOpts.to) { + timeRangeHint= `Until ${todayAwareFormat(dayjs.unix(currOpts.to))}`; + } else if(currOpts.to) { + timeRangeHint = `From ${todayAwareFormat(dayjs.unix(currOpts.from))}`; + } + while (more) { + requestCount++; + const reqOptsHint: string[] = [ + `Page ${currOpts.page}` + ]; + if(timeRangeHint !== undefined) { + reqOptsHint.push(timeRangeHint); + } + if(currOpts.limit !== undefined) { + reqOptsHint.push(`Limit ${currOpts.limit}`); + } + reqLogger.debug(`Fetching => ${reqOptsHint.join(' | ')}`); + const results = await fetcher.getPaginatedTimeRangeListens(currOpts); + if(initial) { + initial = false; + const initialFetchLog = []; + if(results.meta.total !== undefined) { + initialFetchLog.push(`API reported ${results.meta.total} total results`); + } + if(results.meta.limit !== undefined && results.meta.limit !== currOpts.limit) { + initialFetchLog.push(`API reported new limit ${results.meta.limit}`); + currOpts.limit = results.meta.limit; + } + if(initialFetchLog.length > 0) { + logger.debug(initialFetchLog.join(' | ')); + } + } + reqLogger.debug(`${results.data.length} results returned${results.data.length === 0 ? ', ending fetch' : ''}`); + plays = plays.concat(results.data); + if (!results.meta.more) { + logger.debug('API indicated no more results, ending fetch'); + more = false; + } + if (results.data.length === 0) { + more = false; + } + // failsafe? + if(more && results.meta.limit !== undefined && results.data.length < results.meta.limit) { + reqLogger.debug(`Number of returned results was less than reported/defined limit (${results.meta.limit}), ending fetch`); + more = false; + } + + if(more && opts.to === undefined && opts.from === undefined) { + // only wanted one fetch + logger.debug('No to/from defined, ending fetch'); + more = false; + } + + if(more) { + currOpts.page++; + } + } + return plays; + } + } + + throw new Error('fetcher does not implement pagination interface'); +} \ No newline at end of file -- 2.51.2