From b498939d9e2b0b790510677391d7d5dac2b54797 Mon Sep 17 00:00:00 2001 From: FoxxMD Date: Mon, 17 Aug 2026 16:13:03 +0000 Subject: [PATCH] fix: Add queue id to consumeQueue and implement queued id map to prevent processing an already processing play --- src/backend/sources/AbstractSource.ts | 16 +++++++++++++--- src/backend/utils/AsyncUtils.ts | 19 +++++++++++-------- 2 files changed, 24 insertions(+), 11 deletions(-) diff --git a/src/backend/sources/AbstractSource.ts b/src/backend/sources/AbstractSource.ts index 80fcbf0c..c1ff4508 100644 --- a/src/backend/sources/AbstractSource.ts +++ b/src/backend/sources/AbstractSource.ts @@ -902,10 +902,17 @@ export default abstract class AbstractSource extends AbstractComponent implement } = {}, } = this.config; const maxRetries = Math.max(0, maxRequestRetries); + const consumedIds = new Map(); try { await consumeQueue( - () => this.playRepo.getQueueNext(INGRESS_QUEUE), + async (queueId) => { + const next = await this.playRepo.getQueueNext(INGRESS_QUEUE, {notIds: consumedIds.size === 0 ? undefined : consumedIds.values().toArray()}); + if(next !== undefined) { + consumedIds.set(queueId, next.id); + } + return next; + }, async (item) => { if (taskFailures > 0) { const delayFor = pollingBackoff(taskFailures + 1, retryMultiplier); @@ -918,10 +925,12 @@ export default abstract class AbstractSource extends AbstractComponent implement concurrency: this.queueConcurrency, idleMs: this.queueIdleMs, signal, - onSuccess: () => { + onSuccess: (item, queueId) => { + consumedIds.delete(queueId); taskFailures = Math.max(taskFailures - 1, 0); }, - onError: async (e: Error) => { + onError: async (e: Error, queueId) => { + consumedIds.delete(queueId); taskFailures++; this.emitter.emit('discoveryQueueError', e); if(taskFailures < maxRetries) { @@ -1015,6 +1024,7 @@ export default abstract class AbstractSource extends AbstractComponent implement if(state === 'discovered') { await this.scrobble([{...currQueuedPlay.play, id: currQueuedPlay.id, uid: currQueuedPlay.uid}]); } + return currQueuedPlay; } protected setIsSleeping(sleeping: boolean) { diff --git a/src/backend/utils/AsyncUtils.ts b/src/backend/utils/AsyncUtils.ts index b86a0141..17e673b0 100644 --- a/src/backend/utils/AsyncUtils.ts +++ b/src/backend/utils/AsyncUtils.ts @@ -1,5 +1,6 @@ import type {Mapper} from "p-map"; import { sleep } from "../utils.ts"; +import { nanoid } from "nanoid"; /** https://stackoverflow.com/a/63795192/1469797 */ export async function findAsyncSequential( @@ -105,18 +106,18 @@ export const consumeQueueOnce = async (next: () => Promise, pr }; export const consumeQueue = async ( - next: () => Promise, - process: (item: T) => Promise, + next: (queueId: string) => Promise, + process: (item: T, queueId: string) => Promise, opts: { concurrency: number; idleMs: number; signal: AbortSignal; - onError?: (e: Error) => Promise, - onSuccess?: () => void, + onError?: (e: Error, queueId: string) => Promise, + onSuccess?: (item: T, queueId: string) => void, onEmpty?: () => void }, ): Promise => { - const { concurrency, idleMs, signal, onError, onEmpty } = opts; + const { concurrency, idleMs, signal, onError, onEmpty, onSuccess } = opts; while (true) { signal.throwIfAborted(); const inFlight = new Set>(); @@ -127,13 +128,15 @@ export const consumeQueue = async ( await Promise.race(inFlight); continue; } - const item = await next(); + const qId = nanoid(); + const item = await next(qId); if (item === undefined) break; const task = (async () => { try { - await process(item); + await process(item, qId); + onSuccess?.(item, qId); } catch (err) { - await onError?.(err); // swallow so one bad item doesn't kill the loop + await onError?.(err, qId); // swallow so one bad item doesn't kill the loop } })(); inFlight.add(task); -- 2.51.2