From 6f08837665c73d4f564f956deb383dd1773b1b6e Mon Sep 17 00:00:00 2001 From: FoxxMD Date: Fri, 14 Aug 2026 15:10:03 +0000 Subject: [PATCH] refactor: rename queue constant to remove client will be using for both source and client so name is less restrictive now --- .../common/database/drizzle/entityUtils.ts | 4 +- .../drizzle/repositories/QueueRepository.ts | 6 +-- .../scrobblers/AbstractScrobbleClient.ts | 54 +++++++++---------- src/backend/server/api.ts | 4 +- .../tests/scrobbler/scrobblers.test.ts | 16 +++--- src/core/Atomic.ts | 6 +-- src/core/tests/utils/apiFixtures.ts | 4 +- 7 files changed, 47 insertions(+), 47 deletions(-) diff --git a/src/backend/common/database/drizzle/entityUtils.ts b/src/backend/common/database/drizzle/entityUtils.ts index 5056bc63..76d133e7 100644 --- a/src/backend/common/database/drizzle/entityUtils.ts +++ b/src/backend/common/database/drizzle/entityUtils.ts @@ -4,7 +4,7 @@ import type {PlayInputNew} from "./drizzleTypes.ts"; import type {QueueStateNew} from "./drizzleTypes.ts"; import type {ComponentNew} from "./drizzleTypes.ts"; import type { MarkOptional } from "ts-essentials"; -import { CLIENT_DEAD_QUEUE, type DeadLetterScrobble, type ErrorLike, type PlayObject } from "../../../../core/Atomic.ts"; +import { DEAD_QUEUE, type DeadLetterScrobble, type ErrorLike, type PlayObject } from "../../../../core/Atomic.ts"; import dayjs from "dayjs"; import { playContentBasicInvariantTransform, playMbidIdentifier } from "../../../utils/PlayComparisonUtils.ts"; import { hashObject } from "../../../utils/StringUtils.ts"; @@ -72,7 +72,7 @@ export const hydratePlaySelect = (s } export const playSelectToDeadScrobble = (select: PlaySelectWithQueueStates, serializedError: boolean = false): DeadLetterScrobble => { - const deadQueue = select.queueStates.find(x => x.queueName === CLIENT_DEAD_QUEUE); + const deadQueue = select.queueStates.find(x => x.queueName === DEAD_QUEUE); return { play: select.play, id: select.uid, diff --git a/src/backend/common/database/drizzle/repositories/QueueRepository.ts b/src/backend/common/database/drizzle/repositories/QueueRepository.ts index 8a5a4eba..52e53e78 100644 --- a/src/backend/common/database/drizzle/repositories/QueueRepository.ts +++ b/src/backend/common/database/drizzle/repositories/QueueRepository.ts @@ -3,7 +3,7 @@ import { DrizzleBaseRepository, type DrizzleRepositoryOpts } from "./BaseReposit import type {DbConcrete} from "../drizzleUtils.ts"; import type {QueueStateSelect} from "../drizzleTypes.ts"; import { queueStates } from "../schema/schema.ts"; -import { CLIENT_DEAD_QUEUE } from "../../../../../core/Atomic.ts"; +import { DEAD_QUEUE } from "../../../../../core/Atomic.ts"; export class DrizzleQueueRepository extends DrizzleBaseRepository<'queueStates'> { constructor(db: DbConcrete, opts: DrizzleRepositoryOpts = {}) { @@ -17,7 +17,7 @@ export class DrizzleQueueRepository extends DrizzleBaseRepository<'queueStates'> eq(queueStates.componentId, componentId), lte(queueStates.retries, retries), eq(queueStates.queueStatus, 'failed'), - eq(queueStates.queueName, CLIENT_DEAD_QUEUE) + eq(queueStates.queueName, DEAD_QUEUE) )); } @@ -27,7 +27,7 @@ export class DrizzleQueueRepository extends DrizzleBaseRepository<'queueStates'> }).where(and( eq(queueStates.componentId, componentId), eq(queueStates.queueStatus, 'queued'), - eq(queueStates.queueName, CLIENT_DEAD_QUEUE) + eq(queueStates.queueName, DEAD_QUEUE) )); } diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index 2f317e72..12df1f52 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -10,8 +10,8 @@ import { type PlayObject, type QueuedScrobble, type ScrobbleActionResult, type PlayMatchResult, type SourcePlayerObj, type ErrorLike, - CLIENT_INGRESS_QUEUE, - CLIENT_DEAD_QUEUE, + INGRESS_QUEUE, + DEAD_QUEUE, type PlayOriginal, type PlayLifecycle, type SourcePlayerJson, @@ -167,7 +167,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i this.logger = childLogger(logger, this.getIdentifier()); this.npLogger = childLogger(this.logger, 'Now Playing'); this.dupeLogger = childLogger(this.logger, 'Dupe'); - this.deadLogger = childLogger(this.logger, CLIENT_DEAD_QUEUE); + this.deadLogger = childLogger(this.logger, DEAD_QUEUE); this.emitter = emitter; const { @@ -397,17 +397,17 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i if(scrobbledCount !== undefined) { this.tracksScrobbledTotal = scrobbledCount['count(*)']; } - await this.updateQueueStats([CLIENT_INGRESS_QUEUE, CLIENT_DEAD_QUEUE]); + await this.updateQueueStats([INGRESS_QUEUE, DEAD_QUEUE]); } protected async updateQueueStats(queueNames: string[]) { - if(queueNames.includes(CLIENT_INGRESS_QUEUE)) { - this.queuedLength = await this.queueRepo.getQueueCount(this.dbComponent.id, [CLIENT_INGRESS_QUEUE]); + if(queueNames.includes(INGRESS_QUEUE)) { + this.queuedLength = await this.queueRepo.getQueueCount(this.dbComponent.id, [INGRESS_QUEUE]); this.queuedGauge.labels(this.getPrometheusLabels()).set(this.queuedLength); } - if(queueNames.includes(CLIENT_DEAD_QUEUE)) { - this.deadLetterLength = await this.queueRepo.getQueueCount(this.dbComponent.id, [CLIENT_DEAD_QUEUE], ['queued', 'failed']); - this.deadLetterQueued = await this.queueRepo.getQueueCount(this.dbComponent.id, [CLIENT_DEAD_QUEUE], ['queued']); + if(queueNames.includes(DEAD_QUEUE)) { + this.deadLetterLength = await this.queueRepo.getQueueCount(this.dbComponent.id, [DEAD_QUEUE], ['queued', 'failed']); + this.deadLetterQueued = await this.queueRepo.getQueueCount(this.dbComponent.id, [DEAD_QUEUE], ['queued']); // TODO this.deadLetterGauge.labels(this.getPrometheusLabels()).set(this.deadLetterLength); } @@ -883,8 +883,8 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i } handleQueuedScrobbleRanges = async (deadRetries: number = 3) => { - const queued = await this.playRepo.getQueuedScrobbleRange(CLIENT_INGRESS_QUEUE); - const dead = await this.playRepo.getQueuedScrobbleRange(CLIENT_DEAD_QUEUE, {retries: deadRetries}); + const queued = await this.playRepo.getQueuedScrobbleRange(INGRESS_QUEUE); + const dead = await this.playRepo.getQueuedScrobbleRange(DEAD_QUEUE, {retries: deadRetries}); this.scrobbleSOTRanges = groupPlaysToTimeRanges(queued.concat(dead), this.scrobbleSOTRanges, {staleNowBuffer: this.config.options?.refreshStaleAfter}); } @@ -1179,7 +1179,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i while (true) { signal.throwIfAborted(); //let queueEmpty = await this.playRepo.hasQueueNext(CLIENT_INGRESS_QUEUE); // this.queuedLength; // this.queuedScrobbles.length === 0; - let nextQueued = await this.playRepo.getQueueNext(CLIENT_INGRESS_QUEUE); + let nextQueued = await this.playRepo.getQueueNext(INGRESS_QUEUE); if(nextQueued !== undefined) { while (nextQueued !== undefined) { await this.processQueueCurrentScrobble(nextQueued, signal); @@ -1188,7 +1188,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i this.errors = []; this.emitComponentUpdate>({errors: []}); } - nextQueued = await this.playRepo.getQueueNext(CLIENT_INGRESS_QUEUE) + nextQueued = await this.playRepo.getQueueNext(INGRESS_QUEUE) } this.emitEvent('queueEmptied', {}); this.setStatus('Waiting for Plays from Sources'); @@ -1324,7 +1324,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i // } throw e; } finally { - const queueState = currQueuedPlay.queueStates.find(x => x.queueName === CLIENT_INGRESS_QUEUE); + const queueState = currQueuedPlay.queueStates.find(x => x.queueName === INGRESS_QUEUE); if(queueError !== undefined) { await this.queueRepo.updateById(queueState.id, { queueStatus: 'failed', @@ -1385,10 +1385,10 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i this.emitEvent('queueState', {queueName: 'dead', status: 'Running'}); await this.queueRepo.deadFailedToQueue(this.dbComponent.id, retries); - const processable = await this.queueRepo.getQueueCount(this.dbComponent.id, [CLIENT_DEAD_QUEUE]); //this.deadLetterScrobbles.filter(x => x.retries < retries); + const processable = await this.queueRepo.getQueueCount(this.dbComponent.id, [DEAD_QUEUE]); //this.deadLetterScrobbles.filter(x => x.retries < retries); this.deadLetterQueued = processable; - const total = await this.queueRepo.getQueueCount(this.dbComponent.id, [CLIENT_DEAD_QUEUE], ['queued','failed']); + const total = await this.queueRepo.getQueueCount(this.dbComponent.id, [DEAD_QUEUE], ['queued','failed']); this.deadLetterLength = total; const queueStatus = `${processable} of ${total} dead scrobbles have less than ${retries} retries, ${processable === 0 ? 'will skip processing.': 'processing now...'}`; if (processable === 0) { @@ -1402,7 +1402,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i } // await this.handleQueuedScrobbleRanges(); - let nextQueued: PlaySelectWithQueueStates = await this.playRepo.getQueueNext(CLIENT_DEAD_QUEUE, {retries}); + let nextQueued: PlaySelectWithQueueStates = await this.playRepo.getQueueNext(DEAD_QUEUE, {retries}); if(nextQueued !== undefined) { while(nextQueued !== undefined) { const [scrobbled, dead] = await this.processDeadLetterScrobble(nextQueued.uid, signal); @@ -1410,7 +1410,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i if(scrobbled) { removedIds.push(dead.id); } - nextQueued = await this.playRepo.getQueueNext(CLIENT_DEAD_QUEUE, {retries}); + nextQueued = await this.playRepo.getQueueNext(DEAD_QUEUE, {retries}); } } @@ -1446,7 +1446,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i if(deadScrobble.state === 'scrobbled') { throw new Error(`Play ${uid} is already scrobbled.`); } - const deadQueueState: QueueStateSelect = deadScrobble.queueStates.find(x => x.queueName === CLIENT_DEAD_QUEUE); + const deadQueueState: QueueStateSelect = deadScrobble.queueStates.find(x => x.queueName === DEAD_QUEUE); if(deadQueueState === undefined) { throw new Error(`Play ${uid} is not currently queued in dead letter.`); } @@ -1554,7 +1554,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i // this.deadLogger.warn(`No scrobble found with ID ${id}`); // return; // } - const deadQueueState = deadScrobble.queueStates.find(x => x.queueName === CLIENT_DEAD_QUEUE && x.queueStatus !== 'completed'); + const deadQueueState = deadScrobble.queueStates.find(x => x.queueName === DEAD_QUEUE && x.queueStatus !== 'completed'); if(deadQueueState === undefined) { throw new Error(`Play ${deadScrobble.uid} is not currently queued in dead letter.`); } @@ -1585,7 +1585,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i const ids = await this.playRepo.findPlayIdentifiers({ queues: [ { - queueName: CLIENT_DEAD_QUEUE, + queueName: DEAD_QUEUE, queueStatus: types } ] @@ -1593,7 +1593,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i this.deadLogger.info(`Marking ${ids} as completed but unsuccessful...`); await Promise.all(ids.map((x) => this.removeDeadLetterScrobble(x, state, success))); this.deadLogger.info('Finished processing dead scrobbles.'); - await this.updateQueueStats([CLIENT_DEAD_QUEUE]); + await this.updateQueueStats([DEAD_QUEUE]); } queueScrobble = async (data: PlayObject | PlayObject[], source: string, transformFunc?: (x: PlayObject) => Promise) => { @@ -1604,9 +1604,9 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i for await(const play of pMapIterable(playDatas, this.staggerMappers.preCompare(async x => transformFunc !== undefined ? await transformFunc(x) : await this.transformPlay(x, TRANSFORM_HOOK.preCompare)), {concurrency: 3})) { try { // cheap check, looks for play data (non-meta) hash, playdate, and optionally mbid recording - const cheapExisting = await this.playRepo.checkExisting(play, { queueName: CLIENT_INGRESS_QUEUE }); + const cheapExisting = await this.playRepo.checkExisting(play, { queueName: INGRESS_QUEUE }); if (cheapExisting !== undefined) { - const qs = cheapExisting.queueStates.find(x => x.queueName === CLIENT_INGRESS_QUEUE); + const qs = cheapExisting.queueStates.find(x => x.queueName === INGRESS_QUEUE); this.logger.trace(`Not adding to queue because it is already in the queue, discovered via hash/mbid, last queued at ${todayAwareFormat(qs.createdAt)}`); continue; } @@ -1614,7 +1614,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i let offset = 0; let inQueue = false; while (true) { - const { data, meta } = await this.playRepo.getQueued(CLIENT_INGRESS_QUEUE, { offset }); + const { data, meta } = await this.playRepo.getQueued(INGRESS_QUEUE, { offset }); const existingQueued = await this.existingScrobble(play, data.map(x => asPlay(x.play)), false); // want to be very confident of this if (existingQueued.match && existingQueued.score > 0.99) { @@ -1650,7 +1650,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i }); const playRow = await this.playRepo.createPlays([createPlayData]); - const queueState = await this.queueRepo.create({componentId: this.dbComponent.id, playId: playRow[0].id, queueName: CLIENT_INGRESS_QUEUE}); + const queueState = await this.queueRepo.create({componentId: this.dbComponent.id, playId: playRow[0].id, queueName: INGRESS_QUEUE}); createdQueuedPlays.push(playRow[0]); this.logger.debug(`Added ${buildTrackString(play)} to the queue`); this.setStatus(`Added Play from parent ${play.uid} to queue`); @@ -1682,7 +1682,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i const newQueue = await this.queueRepo.create({ componentId: this.dbComponent.id, playId: data.id, - queueName: CLIENT_DEAD_QUEUE + queueName: DEAD_QUEUE }) as QueueStateSelect; const deadData = {id: nanoid(), retries: 0, error: e, play: data.play}; //this.deadLetterScrobbles.push(deadData); diff --git a/src/backend/server/api.ts b/src/backend/server/api.ts index 4fdbebdd..474f2b0a 100644 --- a/src/backend/server/api.ts +++ b/src/backend/server/api.ts @@ -6,7 +6,7 @@ import { FixedSizeList } from 'fixed-size-list'; import type { PassThrough } from "node:stream"; import { Transform } from "stream"; import { - CLIENT_DEAD_QUEUE, + DEAD_QUEUE, type ClientStatusData, type DeadLetterScrobble, type LogOutputConfig, @@ -599,7 +599,7 @@ export const setupApi = (app: Express, logger: Logger, appLoggerStream: PassThro ...query as Partial, queues: [ { - queueName: CLIENT_DEAD_QUEUE, + queueName: DEAD_QUEUE, queueStatus: ['queued','failed'] } ] diff --git a/src/backend/tests/scrobbler/scrobblers.test.ts b/src/backend/tests/scrobbler/scrobblers.test.ts index c801d37b..a402b1d8 100644 --- a/src/backend/tests/scrobbler/scrobblers.test.ts +++ b/src/backend/tests/scrobbler/scrobblers.test.ts @@ -6,7 +6,7 @@ import dayjs from "dayjs"; import { after, describe, it } from 'mocha'; import { http, HttpResponse } from 'msw'; import pEvent from 'p-event'; -import { CLIENT_INGRESS_QUEUE, type PlayObject, SOURCE_SOT } from "../../../core/Atomic.ts"; +import { INGRESS_QUEUE, type PlayObject, SOURCE_SOT } from "../../../core/Atomic.ts"; import { sleep, sortByOldestPlayDate } from "../../utils.ts"; import { genGroupIdStr } from '../../../core/PlayUtils.ts'; import mixedDuration from '../plays/mixedDuration.json' with { type: 'json' }; @@ -678,7 +678,7 @@ describe('Scrobble client uses transform plays correctly', function() { track: 'my cool track' }); await testScrobbler.queueScrobble(newScrobble, 'test'); - const queuedPlayedData = await testScrobbler.playRepoTest.getQueued(CLIENT_INGRESS_QUEUE); + const queuedPlayedData = await testScrobbler.playRepoTest.getQueued(INGRESS_QUEUE); expect(queuedPlayedData.data[0].play.data.track).is.eq('my cool track'); testScrobbler.scrobbleSleep = 100; testScrobbler.initScrobbleMonitoring().catch(console.error); @@ -1187,7 +1187,7 @@ describe('Scrobble Clients Behavior', function() { pEvent(testClient.emitter, 'scrobbleQueued'), sleep(10) ]); - const queued = await testClient.getQueued(CLIENT_INGRESS_QUEUE); + const queued = await testClient.getQueued(INGRESS_QUEUE); expect(queued.data).is.empty; }); @@ -1214,7 +1214,7 @@ describe('Scrobble Clients Behavior', function() { pEvent(testClient.emitter, 'scrobbleQueued'), sleep(100) ]) - const queued = await testClient.getQueued(CLIENT_INGRESS_QUEUE); + const queued = await testClient.getQueued(INGRESS_QUEUE); expect(queued.data).is.not.empty; }); @@ -1241,7 +1241,7 @@ describe('Scrobble Clients Behavior', function() { pEvent(testClient.emitter, 'scrobbleQueued'), sleep(50) ]) - const queued = await testClient.getQueued(CLIENT_INGRESS_QUEUE); + const queued = await testClient.getQueued(INGRESS_QUEUE); expect(queued.data).is.not.empty; }); @@ -1268,7 +1268,7 @@ describe('Scrobble Clients Behavior', function() { pEvent(testClient.emitter, 'scrobbleQueued'), sleep(50) ]) - const queued = await testClient.getQueued(CLIENT_INGRESS_QUEUE); + const queued = await testClient.getQueued(INGRESS_QUEUE); expect(queued.data).is.not.empty; }); @@ -1295,7 +1295,7 @@ describe('Scrobble Clients Behavior', function() { pEvent(testClient.emitter, 'scrobbleQueued'), sleep(50) ]) - const queued = await testClient.getQueued(CLIENT_INGRESS_QUEUE); + const queued = await testClient.getQueued(INGRESS_QUEUE); expect(queued.data).is.not.empty; }); @@ -1322,7 +1322,7 @@ describe('Scrobble Clients Behavior', function() { pEvent(testClient.emitter, 'scrobbleQueued'), sleep(50) ]) - const queued = await testClient.getQueued(CLIENT_INGRESS_QUEUE); + const queued = await testClient.getQueued(INGRESS_QUEUE); expect(queued.data).is.not.empty; }); diff --git a/src/core/Atomic.ts b/src/core/Atomic.ts index c10f76e9..b3cf4c66 100644 --- a/src/core/Atomic.ts +++ b/src/core/Atomic.ts @@ -607,10 +607,10 @@ export const REGEX_ISO8601_LOOSE = new RegExp(/\d{4}-[01]\d-[0-3]\dT/); */ export const REGEX_ISO8601_WELLKNOWN = new RegExp(/dayjs-(\d{4}-[01]\d-[0-3]\dT.*)/); -export const CLIENT_INGRESS_QUEUE: QueueName = 'ingress'; -export const CLIENT_DEAD_QUEUE: QueueName = 'dead'; +export const INGRESS_QUEUE: QueueName = 'ingress'; +export const DEAD_QUEUE: QueueName = 'dead'; export type QueueName = 'ingress' | 'dead'; -export const QUEUE_NAMES = [CLIENT_INGRESS_QUEUE, CLIENT_DEAD_QUEUE]; +export const QUEUE_NAMES = [INGRESS_QUEUE, DEAD_QUEUE]; /** * Useful TS type-only utility for testing type equality diff --git a/src/core/tests/utils/apiFixtures.ts b/src/core/tests/utils/apiFixtures.ts index d5f07d66..9d3730f3 100644 --- a/src/core/tests/utils/apiFixtures.ts +++ b/src/core/tests/utils/apiFixtures.ts @@ -1,6 +1,6 @@ import { faker } from "@faker-js/faker"; import type {ComponentClientApi, ComponentClientApiJson, ComponentCommonApi, ComponentCommonApiJson, ComponentSourceApi, ComponentSourceApiJson, ComponentState, PlayApiCommon, PlayApiCommonDetailed, PlayInputApi, QueueStateApi} from "../../Api.ts"; -import { CLIENT_INGRESS_QUEUE, COMPONENT_AUTH_TYPE, type ComponentType, type JsonPlayObject, type PlayObject, QUEUE_STATUSES, type SourcePlayerJson, sourceSotTypes } from "../../Atomic.ts"; +import { INGRESS_QUEUE, COMPONENT_AUTH_TYPE, type ComponentType, type JsonPlayObject, type PlayObject, QUEUE_STATUSES, type SourcePlayerJson, sourceSotTypes } from "../../Atomic.ts"; import { generatePlay, normalizePlays } from "./PlayTestUtils.ts"; import { generatePlayInput, generatePlayWithLifecycle, playWithLifecycleScrobble, randomPlayState } from "./fixtures.ts"; import { asJsonPlayObject } from "../../PlayMarshalUtils.ts"; @@ -73,7 +73,7 @@ export const generateQueueStateApi = (data: Partial): QueueStateA const cAt = faker.date.recent().toISOString(); return { id: faker.number.int({min: 1, max: 100}), - queueName: CLIENT_INGRESS_QUEUE, + queueName: INGRESS_QUEUE, queueStatus: faker.helpers.arrayElement(QUEUE_STATUSES), updatedAt: cAt, retries: 0, -- 2.51.2