diff --git a/src/backend/common/database/drizzle/repositories/PlayRepository.ts b/src/backend/common/database/drizzle/repositories/PlayRepository.ts index b84767e7..3098576d 100644 --- a/src/backend/common/database/drizzle/repositories/PlayRepository.ts +++ b/src/backend/common/database/drizzle/repositories/PlayRepository.ts @@ -59,6 +59,7 @@ export class DrizzlePlayRepository extends DrizzleBaseRepository<'plays'> { protected hasQueueNextPrepared?: ReturnType protected getQueueNextPrepared?: ReturnType + protected getQueuedScrobbleRangePrepared?: ReturnType constructor(db: ReturnType, opts: DrizzleRepositoryOpts = {}) { super(db, 'plays', 'Plays', opts); @@ -487,6 +488,31 @@ export class DrizzlePlayRepository extends DrizzleBaseRepository<'plays'> { return nextId !== undefined; } + protected prepareGetQueuedScrobbleRange = () => this.db.query.plays.findMany({ + where: { + componentId: this.componentId, + queueStates: { + queueName: sql.placeholder('queueName'), + queueStatus: 'queued', + retries: { + lte: sql.placeholder('retries') + } + }, + }, + orderBy: { + seenAt: 'asc', + }, + limit: sql.placeholder('limit') + }).prepare() + + public getQueuedScrobbleRange = async (queueName: string, opts: {retries?: number, limit?: number} = {}): Promise => { + if(this.getQueuedScrobbleRangePrepared === undefined) { + this.getQueuedScrobbleRangePrepared = this.prepareGetQueuedScrobbleRange(); + } + const res = await this.getQueuedScrobbleRangePrepared.execute({queueName, retries: opts.retries ?? 0, limit: opts.limit ?? 30}); + return res.map(x => x.play); + } + public getQueued = async (queueName: string, opts: { order?: 'asc' | 'desc', limit?: number, diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index 66702caa..8ff8344e 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -587,8 +587,8 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i } handleQueuedScrobbleRanges = async (deadRetries: number = 3) => { - const queued = (await this.playRepo.getQueued(CLIENT_INGRESS_QUEUE, {limit: 30})).data.map(x => asPlay(x.play)); - const dead = (await this.playRepo.getQueued(CLIENT_DEAD_QUEUE, {limit: 30, retries: deadRetries})).data.map(x => asPlay(x.play)); + const queued = await this.playRepo.getQueuedScrobbleRange(CLIENT_INGRESS_QUEUE); + const dead = await this.playRepo.getQueuedScrobbleRange(CLIENT_DEAD_QUEUE, {retries: deadRetries}); this.scrobbleSOTRanges = groupPlaysToTimeRanges(queued.concat(dead), this.scrobbleSOTRanges, {staleNowBuffer: this.config.options?.refreshStaleAfter}); }