From 68e473ec0286aaabc7165c3701ca6fc07896179c Mon Sep 17 00:00:00 2001 From: FoxxMD Date: Tue, 24 Mar 2026 20:21:14 +0000 Subject: [PATCH] fix: Add stagger options as injected dependency to make tests faster --- src/backend/ioc.ts | 6 +++++- .../scrobblers/AbstractScrobbleClient.ts | 19 +++++++++++++------ src/backend/sources/AbstractSource.ts | 9 ++++++--- src/backend/tests/scrobbler/TestScrobbler.ts | 4 ++-- src/backend/tests/setup.ts | 2 +- src/backend/tests/utils/CacheTestUtils.ts | 2 +- 6 files changed, 28 insertions(+), 14 deletions(-) diff --git a/src/backend/ioc.ts b/src/backend/ioc.ts index 1cd037c5..2df8c8aa 100644 --- a/src/backend/ioc.ts +++ b/src/backend/ioc.ts @@ -14,6 +14,7 @@ import { TransformerCommonConfig } from "../core/Atomic.js"; import prom, { Counter, Gauge } from 'prom-client'; import { CoverArtApiClient } from "./common/vendor/musicbrainz/CoverArtApiClient.js"; import { version } from "./version.js"; +import { StaggerOptions } from "./utils/AsyncUtils.js"; let root: ReturnType; export interface RootOptions { @@ -26,6 +27,7 @@ export interface RootOptions { cache?: CacheConfigOptions | MSCache | (() => MSCache) mbMap?: MusicBrainzSingletonMap | (() => MusicBrainzSingletonMap) transformers?: TransformerCommonConfig[] + staggerOptions?: Partial } const discovered = new prom.Counter({ @@ -60,6 +62,7 @@ const createRoot = (options: RootOptions = {logger: loggerDebug}) => { cache, mbMap, transformers = [], + staggerOptions, } = options || {}; const configDir = process.env.CONFIG_DIR || path.resolve(projectDir, `./config`); let disableWeb = dw; @@ -146,7 +149,8 @@ const createRoot = (options: RootOptions = {logger: loggerDebug}) => { transformerManager, cache: () => maybeSingletonCache !== undefined ? () => maybeSingletonCache : cacheFunc, mbMap: () => maybeSingletonMb !== undefined ? () => maybeSingletonMb : mbFunc, - coverArtApi + coverArtApi, + staggerOptions: staggerOptions ?? {}, }).add((items) => { const localUrl = generateBaseURL(baseUrl, items.port) return { diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index 6d896bb9..542ae800 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -58,7 +58,7 @@ import { WebhookPayload } from "../common/infrastructure/config/health/webhooks. import { AsyncTask, SimpleIntervalJob, Task, ToadScheduler } from "toad-scheduler"; import { getRoot } from "../ioc.js"; import { rehydratePlay } from "../utils/CacheUtils.js"; -import { findAsyncSequential, staggerMapper } from "../utils/AsyncUtils.js"; +import { findAsyncSequential, staggerMapper, StaggerOptions } from "../utils/AsyncUtils.js"; import pMap, { pMapIterable } from "p-map"; import { comparePlayArtistsNormalized, comparePlayTracksNormalized, existingScrobble, ExistingScrobbleOpts } from "../utils/PlayComparisonUtils.js"; import { lifecyclelessInvariantTransform } from "../../core/PlayUtils.js"; @@ -124,6 +124,8 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i protected deadLetterGauge: Gauge; protected problemGauge: Gauge; + protected staggerOpts: Partial; + constructor(type: any, name: any, config: CommonClientConfig, notifier: Notifiers, emitter: EventEmitter, logger: Logger) { super(config); this.type = type; @@ -187,7 +189,8 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i transformPlay: this.transformPlay, existingSubmitted: this.findExistingSubmittedPlayObj } - this.existingScrobble = (playObjPre: PlayObject, existingScrobbles: PlayObject[], log?: boolean) => existingScrobble(playObjPre, existingScrobbles, existingScrobbleOpts, log) + this.existingScrobble = (playObjPre: PlayObject, existingScrobbles: PlayObject[], log?: boolean) => existingScrobble(playObjPre, existingScrobbles, existingScrobbleOpts, log); + this.staggerOpts = getRoot().items.staggerOptions; } protected getIdentifier() { @@ -446,7 +449,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i const playObj = await this.transformPlay(playObjPre, TRANSFORM_HOOK.candidate); - const sm = staggerMapper({concurrency: 2}); + const sm = staggerMapper({...this.staggerOpts, concurrency: 2}); const dtInvariantMatches = (await pMap(this.scrobbledPlayObjs.data, sm(async x => ({...x, play: await this.transformPlay(x.play, TRANSFORM_HOOK.existing)})), {concurrency: 2})) .filter(x => playObjDataMatch(playObj, x.play)); @@ -848,7 +851,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i queueScrobble = async (data: PlayObject | PlayObject[], source: string) => { const plays = (Array.isArray(data) ? data : [data]).map(x => ({...x, meta: {...x.meta, seenAt: dayjs()}})); - const sm = staggerMapper({concurrency: 2}); + const sm = staggerMapper({...this.staggerOpts, concurrency: 2}); for await(const play of pMapIterable(plays, sm(async x => await this.transformPlay(x, TRANSFORM_HOOK.preCompare)), {concurrency: 2})) { try { const existingQueued = await this.existingScrobble(play, this.queuedScrobbles.map(x => x.play), false); @@ -1004,8 +1007,12 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i protected updateQueuedScrobblesCache = () => { this.cache.cacheScrobble.set(`${this.getMachineId()}-queue`, this.queuedScrobbles) - .then(() => null) - .catch((e) => this.logger.warn(new Error('Error while updating queued scrobble cache', {cause: e}))); + .then(() => { + return undefined; + }) + .catch((e) => { + this.logger.warn(new Error('Error while updating queued scrobble cache', {cause: e})) + }); } } diff --git a/src/backend/sources/AbstractSource.ts b/src/backend/sources/AbstractSource.ts index e7ff9783..d2ca2218 100644 --- a/src/backend/sources/AbstractSource.ts +++ b/src/backend/sources/AbstractSource.ts @@ -40,7 +40,7 @@ import { componentFileLogger } from '../common/logging.js'; import { WebhookPayload } from '../common/infrastructure/config/health/webhooks.js'; import { messageWithCauses, messageWithCausesTruncatedDefault } from '../utils/ErrorUtils.js'; import { genericSourcePlayMatch } from '../utils/PlayComparisonUtils.js'; -import { findAsync, staggerMapper } from '../utils/AsyncUtils.js'; +import { findAsync, staggerMapper, StaggerOptions } from '../utils/AsyncUtils.js'; import pMap, {pMapIterable} from 'p-map'; import prom, { Counter, Gauge } from 'prom-client'; import { normalizeStr } from '../utils/StringUtils.js'; @@ -94,6 +94,8 @@ export default abstract class AbstractSource extends AbstractComponent implement protected discoveredCounter: Counter; + protected staggerOpts: Partial; + constructor(type: SourceType, name: string, config: SourceConfig, internal: InternalConfig, emitter: EventEmitter) { super(config); const {clients = [] } = config; @@ -110,6 +112,7 @@ export default abstract class AbstractSource extends AbstractComponent implement this.emitter = emitter; this.discoveredCounter = getRoot().items.sourceMetics.discovered; + this.staggerOpts = getRoot().items.staggerOptions; } protected getIdentifier() { @@ -219,7 +222,7 @@ export default abstract class AbstractSource extends AbstractComponent implement discover = async (plays: PlayObject[], options: { checkAll?: boolean, [key: string]: any } = {}): Promise => { const newDiscoveredPlays: PlayObject[] = []; - const sm = staggerMapper({concurrency: 2}); + const sm = staggerMapper({...this.staggerOpts, concurrency: 2}); for await(const play of pMapIterable(plays, sm(async x => await this.transformPlay(x, TRANSFORM_HOOK.preCompare)), {concurrency: 2})) { if(!(await this.alreadyDiscovered(play, options))) { this.addPlayToDiscovered(play); @@ -251,7 +254,7 @@ export default abstract class AbstractSource extends AbstractComponent implement return; } newDiscoveredPlays.sort(sortByOldestPlayDate); - const sm = staggerMapper({concurrency: 2}); + const sm = staggerMapper({...this.staggerOpts, concurrency: 2}); this.emitter.emit('discoveredToScrobble', { data: await pMap(newDiscoveredPlays, sm(async (x) => await this.transformPlay(x, TRANSFORM_HOOK.postCompare)), {concurrency: 2}), options: { diff --git a/src/backend/tests/scrobbler/TestScrobbler.ts b/src/backend/tests/scrobbler/TestScrobbler.ts index d7847c57..dc4313df 100644 --- a/src/backend/tests/scrobbler/TestScrobbler.ts +++ b/src/backend/tests/scrobbler/TestScrobbler.ts @@ -1,4 +1,3 @@ -import { loggerTest } from "@foxxmd/logging"; import EventEmitter from "events"; import request from "superagent"; import { PlayObject } from "../../../core/Atomic.js"; @@ -7,6 +6,7 @@ import AbstractScrobbleClient from "../../scrobblers/AbstractScrobbleClient.js"; import { CommonClientConfig, CommonClientOptions, NowPlayingOptions } from "../../common/infrastructure/config/client/index.js"; import clone from "clone"; import { TimeRangeListensFetcher } from "../../common/infrastructure/Atomic.js"; +import { loggerNoop } from "../../common/MaybeLogger.js"; export class TestScrobbler extends AbstractScrobbleClient { @@ -14,7 +14,7 @@ export class TestScrobbler extends AbstractScrobbleClient { getScrobblesForTimeRange: TimeRangeListensFetcher; constructor(config: CommonClientConfig = {name: 'test'}) { - const logger = loggerTest; + const logger = loggerNoop; const notifier = new Notifiers(new EventEmitter(), new EventEmitter(), new EventEmitter(), logger); super('test', 'Test', {name: 'test', ...config}, notifier, new EventEmitter(), logger); this.supportsNowPlaying = false; diff --git a/src/backend/tests/setup.ts b/src/backend/tests/setup.ts index ecc3656a..4f143853 100644 --- a/src/backend/tests/setup.ts +++ b/src/backend/tests/setup.ts @@ -2,5 +2,5 @@ import { loggerTest } from '@foxxmd/logging'; import { getRoot } from "../ioc.js"; import { transientCache } from './utils/CacheTestUtils.js'; -const root = getRoot({cache: transientCache, logger: loggerTest}); +const root = getRoot({cache: transientCache, logger: loggerTest, staggerOptions: {initialInterval: 1, maxRandomStagger: 1}}); root.items.cache().init(); \ No newline at end of file diff --git a/src/backend/tests/utils/CacheTestUtils.ts b/src/backend/tests/utils/CacheTestUtils.ts index 96b4621b..a8784271 100644 --- a/src/backend/tests/utils/CacheTestUtils.ts +++ b/src/backend/tests/utils/CacheTestUtils.ts @@ -1,4 +1,4 @@ import { loggerTest } from "@foxxmd/logging"; import { MSCache } from "../../common/Cache.js"; -export const transientCache = () => new MSCache(loggerTest, {scrobble: {provider: 'memory'}}); \ No newline at end of file +export const transientCache = () => new MSCache(loggerTest, {scrobble: {provider: 'memory'}, auth: {provider: 'memory'}}); \ No newline at end of file -- 2.51.2