diff --git a/src/backend/common/AbstractComponent.ts b/src/backend/common/AbstractComponent.ts index a8df7ec4..1f61d147 100644 --- a/src/backend/common/AbstractComponent.ts +++ b/src/backend/common/AbstractComponent.ts @@ -80,12 +80,12 @@ export default abstract class AbstractComponent extends AbstractInitializable { } protected async postInitialize(): Promise { - super.postInitialize(); + await super.postInitialize(); this.componentRepo.updateById(this.dbComponent.id, {lastReadyAt: dayjs()}); } protected async doBuildDatabase(): Promise { - super.doBuildDatabase(); + await super.doBuildDatabase(); let name: string; if('name' in this) { diff --git a/src/backend/common/database/drizzle/repositories/PlayHistoricalRepository.ts b/src/backend/common/database/drizzle/repositories/PlayHistoricalRepository.ts index fef147f5..56771aa2 100644 --- a/src/backend/common/database/drizzle/repositories/PlayHistoricalRepository.ts +++ b/src/backend/common/database/drizzle/repositories/PlayHistoricalRepository.ts @@ -265,11 +265,10 @@ export class DrizzlePlayHistoricalRepository extends DrizzleBaseRepository<'play }) } - public getTemporallyClosePlays = async (play: PlayObject, opts: {states?: PlaySelect['state'][], bufferTime?: number} & ComponentConstrainedRepoOpts = {}): Promise => { + public getTemporallyClosePlays = async (play: PlayObject, opts: {bufferTime?: number} & ComponentConstrainedRepoOpts = {}): Promise => { const { componentId = this.componentId, - bufferTime, - states + bufferTime } = opts; let query: FindMany<'plays'> = {}; @@ -278,11 +277,6 @@ export class DrizzlePlayHistoricalRepository extends DrizzleBaseRepository<'play componentId, playedAt: buildDateCompare(getTemporallyCloseDateCompareOp(play, {bufferTime})), }; - if(states !== undefined) { - where.state = { - in: states - } - } query.where = where; return ((await this.db.query.plays.findMany({ diff --git a/src/backend/scrobblers/AbstractHistoricalScrobbleClient.ts b/src/backend/scrobblers/AbstractHistoricalScrobbleClient.ts index 02df02e3..4f342191 100644 --- a/src/backend/scrobblers/AbstractHistoricalScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractHistoricalScrobbleClient.ts @@ -3,9 +3,12 @@ import { sortByNewestDate } from "../../core/PlayUtils.js"; import AbstractScrobbleClient from "./AbstractScrobbleClient.js"; import { ComponentMigrationSelect } from "../common/database/drizzle/drizzleTypes.js"; import { ErrorIsh } from "../../core/ErrorUtils.js"; -import { DrizzlePlayHistoricalRepository } from "../common/database/drizzle/repositories/PlayHistoricalRepository.js"; +import { DrizzlePlayHistoricalRepository, playToRepositoryCreatePlayHistoricalOpts, RepositoryCreatePlayHistoricalOpts } from "../common/database/drizzle/repositories/PlayHistoricalRepository.js"; import { spawn, isAbortError } from 'abort-controller-x'; import { generateLoggableAbortReason } from "../common/errors/MSErrors.js"; +import { Logger } from "@foxxmd/logging"; +import { buildTrackString } from "../../core/StringUtils.js"; +import { PlayObject } from "../../core/Atomic.js"; export default abstract class AbstractHistoricalScrobbleClient extends AbstractScrobbleClient { @@ -17,6 +20,7 @@ export default abstract class AbstractHistoricalScrobbleClient extends AbstractS synced: boolean; syncedReason?: string; syncError?: ErrorIsh; + override preloadScrobbles: boolean = false; protected abstract doHydrateHistoricalScrobbles(opts: {allowFailures?: boolean, signal?: AbortSignal }): Promise; @@ -75,6 +79,67 @@ export default abstract class AbstractHistoricalScrobbleClient extends AbstractS return [true]; } + protected async createHistoricalPlays(batch: RepositoryCreatePlayHistoricalOpts[], opts: {allowFailures?: boolean, logger?: Logger, signal?: AbortSignal} = {}): Promise<[boolean, number]> { + const { + allowFailures = false, + logger = this.logger, + signal + } = opts; + try { + await this.playsHistoricalRepo.createPlays(batch); + return [true, batch.length]; + } catch (e) { + logger.warn(`Failed to persist batch of ${batch} plays, trying individually...`); + } + signal?.throwIfAborted(); + + let valid = 0; + for(const p of batch) { + try { + await this.playsHistoricalRepo.createPlays([p]); + valid++; + } catch (e) { + if(allowFailures) { + logger.warn(p.play,`Failed to persist play from record with rKey ${p.play.meta.playId} => ${buildTrackString(p.play)}`); + logger.warn(e); + } else { + logger.error(p.play,`Failed to persist play from record with rKey ${p.play.meta.playId} => ${buildTrackString(p.play)}`); + throw e; + } + } + signal?.throwIfAborted(); + } + + return [false, valid]; + } + + async getSOTScrobblesForPlay(play: PlayObject): Promise { + const closeTemporalPlays = await this.playsHistoricalRepo.getTemporallyClosePlays(play); + return closeTemporalPlays.map(x => x.play); + } + + protected abstract syncRecentHistoricalScrobbles(): Promise; + + protected async postInitialize(): Promise { + await super.postInitialize(); + + if(this.lastImport === undefined && this.syncError === undefined) { + // have not run an initial import so automatically do it now + this.logger.info('No historical imports have run! Automatically running an initial import now.'); + this.hydrateHistoricalScrobbles(); + } else { + // pull latest plays into database + this.logger.info('Pulling latest scrobbles to sync up historical database...'); + const recent = await this.syncRecentHistoricalScrobbles(); + if(recent.length > 0) { + await this.createHistoricalPlays(recent.map((x) => playToRepositoryCreatePlayHistoricalOpts({play: x}))); + this.logger.verbose(`Added ${recent.length} upstream plays to historical plays`); + } else { + this.logger.verbose('Most recent plays are already in sync with historical database!'); + } + } + } + protected async postDatabase(): Promise { await super.postDatabase(); this.playsHistoricalRepo = new DrizzlePlayHistoricalRepository(this.db, {componentId: this.dbComponent.id, logger: this.logger}); diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index 3fbfc2b4..e2a583ed 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -99,6 +99,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i protected MAX_STORED_SCROBBLES = 40; protected MAX_INITIAL_SCROBBLES_FETCH = this.MAX_STORED_SCROBBLES; + preloadScrobbles: boolean = true; scrobbleSOTRanges: PaginatedTimeRangeOptions[] = []; tracksScrobbled: number = 0; @@ -714,36 +715,38 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i 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; - } + if(this.preloadScrobbles) { + 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(`Preloading up to ${initialLimit} initial scrobbles...`); + this.logger.verbose(`Preloading up to ${initialLimit} initial scrobbles...`); - try { - const preload = await this.getScrobblesForTimeRange({ - limit: initialLimit, - fetchMax: initialLimit - }); - if(preload === undefined) { - this.logger.warn('Preload result was undefined!'); - } else { - if(preload.length === 0) { - this.logger.verbose(`Preloaded 0 scrobbles.`); + try { + const preload = await this.getScrobblesForTimeRange({ + limit: initialLimit, + fetchMax: initialLimit + }); + if(preload === undefined) { + this.logger.warn('Preload result was undefined!'); } else { - preload.sort(sortByOldestPlayDate); - const from = preload[0].data.playDate; - // we are assuming that all fetchers return latest scrobbles first (pretty sure this is the case) - const to = dayjs();// preload[preload.length - 1].data.playDate; - await this.cache.cacheClientScrobbles.set(this.getScrobbleCacheKey(from, to), preload, '60s'); - this.scrobbleSOTRanges.push({from: from.unix(), to: to.unix()}); - this.logger.verbose(`Preloaded ${preload.length} scrobbles from ${todayAwareFormat(from)} to ${todayAwareFormat(to)}`); + if(preload.length === 0) { + this.logger.verbose(`Preloaded 0 scrobbles.`); + } else { + preload.sort(sortByOldestPlayDate); + const from = preload[0].data.playDate; + // we are assuming that all fetchers return latest scrobbles first (pretty sure this is the case) + const to = dayjs();// preload[preload.length - 1].data.playDate; + await this.cache.cacheClientScrobbles.set(this.getScrobbleCacheKey(from, to), preload, '60s'); + this.scrobbleSOTRanges.push({from: from.unix(), to: to.unix()}); + this.logger.verbose(`Preloaded ${preload.length} scrobbles from ${todayAwareFormat(from)} to ${todayAwareFormat(to)}`); + } } + } catch (e) { + this.logger.warn(new SimpleError('Could not preload scrobbles', {cause: e, shortStack: true})); } - } catch (e) { - this.logger.warn(new SimpleError('Could not preload scrobbles', {cause: e, shortStack: true})); } } @@ -759,7 +762,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i this.scrobbleSOTRanges = groupPlaysToTimeRanges(queued.concat(dead), this.scrobbleSOTRanges, {staleNowBuffer: this.config.options?.refreshStaleAfter}); } - getSOTScrobblesForPlay = async (play: PlayObject): Promise => { + async getSOTScrobblesForPlay(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)}`); @@ -786,6 +789,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i throw new SimpleError('Cannot get historical plays', {cause: e, shortStack: true}); } } + public async alreadyScrobbled(playObj: PlayObject, log?: boolean): Promise<[boolean, PlayMatchResult]> { const result = await this.existingScrobble(playObj, await this.getSOTScrobblesForPlay(playObj)); return [result.match, result]; diff --git a/src/backend/scrobblers/TealfmScrobbler.ts b/src/backend/scrobblers/TealfmScrobbler.ts index 6d9ce995..c94ffc65 100644 --- a/src/backend/scrobblers/TealfmScrobbler.ts +++ b/src/backend/scrobblers/TealfmScrobbler.ts @@ -203,7 +203,8 @@ export default class TealScrobbler extends AbstractHistoricalScrobbleClient { } async fetchCarToFile() { - + // TODO use `since` to get CAR diff instead of entire repo + // can use last import date from migrations table const filename = path.resolve(this.configDir, `${this.getSafeExternalId()}-${dayjs().unix()}.car`); await fsPromise.writeFile(filename, Buffer.from(((await this.client.getCAR())))); return filename; @@ -295,38 +296,16 @@ export default class TealScrobbler extends AbstractHistoricalScrobbleClient { logger.info(`Completed CAR conversion: Result ${allGood ? 'OK' : 'Some Errors'} in ${durationToHuman(dayjs.duration(dayjs().diff(start)))} | Records ${count} | Persisted ${persisted}`) } - async createHistoricalPlays(batch: RepositoryCreatePlayHistoricalOpts[], opts: {allowFailures?: boolean, logger?: Logger, signal?: AbortSignal} = {}): Promise<[boolean, number]> { - const { - allowFailures = false, - logger = this.logger, - signal - } = opts; - try { - await this.playsHistoricalRepo.createPlays(batch); - return [true, batch.length]; - } catch (e) { - logger.warn(`Failed to persist batch of ${batch} plays, trying individually...`); - } - signal?.throwIfAborted(); - - let valid = 0; - for(const p of batch) { - try { - await this.playsHistoricalRepo.createPlays([p]); - valid++; - } catch (e) { - if(allowFailures) { - logger.warn(p.play,`Failed to persist play from record with rKey ${p.play.meta.playId} => ${buildTrackString(p.play)}`); - logger.warn(e); - } else { - logger.error(p.play,`Failed to persist play from record with rKey ${p.play.meta.playId} => ${buildTrackString(p.play)}`); - throw e; - } + protected async syncRecentHistoricalScrobbles(): Promise { + const recentPlays = await this.getScrobblesForTimeRange(undefined); + const unseenPlays: PlayObject[] = []; + for(const p of recentPlays) { + if(!(await this.playsHistoricalRepo.hasByUid(p.meta.playId))) { + unseenPlays.push(p); } - signal?.throwIfAborted(); } - - return [false, valid]; + return unseenPlays; } + }