diff --git a/src/backend/common/AbstractComponent.ts b/src/backend/common/AbstractComponent.ts index 7c940ebc..0661a176 100644 --- a/src/backend/common/AbstractComponent.ts +++ b/src/backend/common/AbstractComponent.ts @@ -24,7 +24,7 @@ import { CommonClientConfig } from "./infrastructure/config/client/index.js"; import { CommonSourceConfig } from "./infrastructure/config/source/index.js"; import play = Simulate.play; import { WebhookPayload } from "./infrastructure/config/health/webhooks.js"; -import { AuthCheckError, BuildDataError, ConnectionCheckError, PostInitError, TransformRulesError } from "./errors/MSErrors.js"; +import { AuthCheckError, BuildDataError, ConnectionCheckError, ParseCacheError, PostInitError, TransformRulesError } from "./errors/MSErrors.js"; import { messageWithCauses, messageWithCausesTruncatedDefault } from "../utils/ErrorUtils.js"; export default abstract class AbstractComponent { @@ -35,6 +35,7 @@ export default abstract class AbstractComponent { buildOK?: boolean | null; connectionOK?: boolean | null; + cacheOK?: boolean | null; initializing: boolean = false; @@ -65,6 +66,7 @@ export default abstract class AbstractComponent { await this.buildComponentLogger(); } await this.buildInitData(force); + await this.parseCache(force); this.buildTransformRules(); await this.checkConnection(force); await this.testAuth(force); @@ -105,6 +107,45 @@ export default abstract class AbstractComponent { } } + public async parseCache(force: boolean = false) { + if(this.cacheOK) { + if(!force) { + return; + } + this.logger.debug('Cache OK but step was forced'); + } + try { + const res = await this.doParseCache(); + if(res === undefined) { + this.cacheOK = null; + this.logger.debug('No cache to parse.'); + return; + } + if (res === true) { + this.logger.verbose('Parsing caching succeeded'); + } else if (typeof res === 'string') { + this.logger.verbose(`Parsing caching succeeded => ${res}`); + } + this.cacheOK = true; + } catch (e) { + this.cacheOK = false; + throw new ParseCacheError('Parsing cache for initialization failed', {cause: e}); + } + } + + /** + * Build or parse any cache required for this Component + * + * * Return undefined if not possible or not required + * * Return TRUE if build succeeded + * * Return string if build succeeded and should log result + * * Throw error on failure + * */ + protected async doParseCache(): Promise { + return; + } + + public async buildInitData(force: boolean = false) { if(this.buildOK) { if(!force) { diff --git a/src/backend/common/Cache.ts b/src/backend/common/Cache.ts index 6d3c4d01..e46c1667 100644 --- a/src/backend/common/Cache.ts +++ b/src/backend/common/Cache.ts @@ -5,32 +5,7 @@ import { childLogger, Logger } from '@foxxmd/logging'; import { projectDir } from './index.js'; import path from 'path'; import { fileOrDirectoryIsWriteable } from '../utils.js'; - -export type CacheProvider = 'memory' | 'valkey' | 'file'; - -interface CacheConfig { - provider: T - connection?: string -} - -export type CacheMetadaProvider = Exclude; -export type CacheMetadataConfig = CacheConfig - -const asCacheMetadataProvider = (val: string): val is CacheScrobbleProvider => { - return ['memory', 'valkey'].includes(val); -} - -export type CacheScrobbleProvider = CacheProvider; -export type CacheScrobbleConfig = CacheConfig; - -const asCacheScrobbleProvider = (val: string): val is CacheScrobbleProvider => { - return ['memory', 'valkey', 'file'].includes(val); -} - -export interface CacheConfigOptions { - metadata?: CacheMetadataConfig - scrobble?: CacheScrobbleConfig -} +import { asCacheMetadataProvider, asCacheScrobbleProvider, CacheConfig, CacheConfigOptions, CacheMetadaProvider, CacheProvider } from './infrastructure/Atomic.js'; const configDir = process.env.CONFIG_DIR || path.resolve(projectDir, `./config`); @@ -43,17 +18,18 @@ export class MSCache { logger: Logger; - constructor(logger: Logger, config: CacheConfigOptions) { + constructor(logger: Logger, config: CacheConfigOptions = {}) { this.logger = childLogger(logger, 'Cache'); const { metadata: { - provider: mProvider = 'memory', + provider: mProvider = (process.env.CACHE_METADATA as (CacheMetadaProvider | undefined) ?? 'memory'), + connection: mConn = process.env.CACHE_METADATA_CONN, ...restMetadata } = {}, scrobble: { - provider: sProvider = 'memory', - connection = configDir, + provider: sProvider = (process.env.CACHE_SCROBBLE as (CacheProvider | undefined) ?? 'file'), + connection = (process.env.CACHE_SCROBBLE_CONN ?? configDir), ...restScrobble } = {}, } = config; @@ -61,6 +37,7 @@ export class MSCache { this.config = { metadata: { provider: mProvider, + connection: mConn, ...restMetadata, }, scrobble: { diff --git a/src/backend/common/errors/MSErrors.ts b/src/backend/common/errors/MSErrors.ts index 3ada4ca2..5060f497 100644 --- a/src/backend/common/errors/MSErrors.ts +++ b/src/backend/common/errors/MSErrors.ts @@ -8,12 +8,16 @@ export class BuildDataError extends StageError { name = 'Init Build Data'; } +export class ParseCacheError extends StageError { + name = 'Init Parse Cache'; +} + export class TransformRulesError extends StageError { name = 'Transform Rules'; } export class ConnectionCheckError extends StageError { - name = 'Conenction Check'; + name = 'Connection Check'; } export class AuthCheckError extends StageError { diff --git a/src/backend/common/infrastructure/Atomic.ts b/src/backend/common/infrastructure/Atomic.ts index 80be609f..2c7f3dc4 100644 --- a/src/backend/common/infrastructure/Atomic.ts +++ b/src/backend/common/infrastructure/Atomic.ts @@ -348,4 +348,24 @@ export type WhenConditionsConfig = WhenConditions; export type WithRequiredProperty = Type & { [Property in Key]-?: Type[Property]; -}; \ No newline at end of file +}; +export type CacheProvider = 'memory' | 'valkey' | 'file'; +export interface CacheConfig { + provider: T; + connection?: string; +} +export type CacheMetadaProvider = Exclude; +export type CacheMetadataConfig = CacheConfig; +export const asCacheMetadataProvider = (val: string): val is CacheScrobbleProvider => { + return ['memory', 'valkey'].includes(val); +}; +export type CacheScrobbleProvider = CacheProvider; +export type CacheScrobbleConfig = CacheConfig; +export const asCacheScrobbleProvider = (val: string): val is CacheScrobbleProvider => { + return ['memory', 'valkey', 'file'].includes(val); +}; +export interface CacheConfigOptions { + metadata?: CacheMetadataConfig; + scrobble?: CacheScrobbleConfig; +} + diff --git a/src/backend/common/infrastructure/config/aioConfig.ts b/src/backend/common/infrastructure/config/aioConfig.ts index cfe6134d..95d5fdce 100644 --- a/src/backend/common/infrastructure/config/aioConfig.ts +++ b/src/backend/common/infrastructure/config/aioConfig.ts @@ -5,7 +5,7 @@ import { RequestRetryOptions } from "./common.js"; import { WebhookConfig } from "./health/webhooks.js"; import { CommonSourceOptions, SourceRetryOptions } from "./source/index.js"; import { SourceAIOConfig } from "./source/sources.js"; -import { ClientType, SourceType } from "../Atomic.js"; +import { CacheConfigOptions, ClientType, SourceType } from "../Atomic.js"; export interface SourceDefaults extends CommonSourceOptions { @@ -64,6 +64,8 @@ export interface AIOConfig { * @examples [false] * */ debugMode?: boolean + + cache?: CacheConfigOptions } export interface AIOClientConfig { diff --git a/src/backend/index.ts b/src/backend/index.ts index 3e47477a..c2a0369d 100644 --- a/src/backend/index.ts +++ b/src/backend/index.ts @@ -103,6 +103,8 @@ const configDir = process.env.CONFIG_DIR || path.resolve(projectDir, `./config`) version: root.get('version') }, root.get('logger')); + await root.get('cache').init(); + initServer(logger, appLoggerStream, output, scrobbleSources, scrobbleClients); if(process.env.IS_LOCAL === 'true') { diff --git a/src/backend/ioc.ts b/src/backend/ioc.ts index ff412106..5cb149f8 100644 --- a/src/backend/ioc.ts +++ b/src/backend/ioc.ts @@ -8,6 +8,8 @@ import { WildcardEmitter } from "./common/WildcardEmitter.js"; import { generateBaseURL } from "./utils/NetworkUtils.js"; import { PassThrough } from "stream"; +import { CacheConfigOptions } from "./common/infrastructure/Atomic.js"; +import { MSCache } from "./common/Cache.js"; export let version: string = 'unknown'; @@ -24,6 +26,7 @@ export interface RootOptions { disableWeb?: boolean loggerStream?: PassThrough loggingConfig?: LogOptions + cache?: CacheConfigOptions } const createRoot = (options?: RootOptions) => { @@ -32,7 +35,7 @@ const createRoot = (options?: RootOptions) => { baseUrl = process.env.BASE_URL, disableWeb: dw, loggerStream, - loggingConfig + loggingConfig, } = options || {}; const configDir = process.env.CONFIG_DIR || path.resolve(projectDir, `./config`); let disableWeb = dw; @@ -62,7 +65,8 @@ const createRoot = (options?: RootOptions) => { notifierEmitter: () => new EventEmitter(), loggerStream, loggingConfig, - logger: options.logger + logger: options.logger, + cache: new MSCache(options.logger, options.cache) }).add((items) => { const localUrl = generateBaseURL(baseUrl, items.port) return { diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index 27ca1c2b..dd77811c 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -52,6 +52,8 @@ import { } from "../utils/TimeUtils.js"; import { WebhookPayload } from "../common/infrastructure/config/health/webhooks.js"; import { AsyncTask, SimpleIntervalJob, Task, ToadScheduler } from "toad-scheduler"; +import { MSCache } from "../common/Cache.js"; +import { getRoot } from "../ioc.js"; type PlatformMappedPlays = Map; type NowPlayingQueue = Map; @@ -97,6 +99,8 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i nowPlayingTaskInterval: number = 5000; npLogger: Logger; + cache: MSCache; + declare config: CommonClientConfig; notifier: Notifiers; @@ -110,6 +114,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i this.npLogger = childLogger(this.logger, 'Now Playing'); this.notifier = notifier; this.emitter = emitter; + this.cache = getRoot().get('cache'); this.scrobbledPlayObjs = new FixedSizeList(this.MAX_STORED_SCROBBLES); @@ -169,6 +174,9 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i protected getIdentifier() { return `${capitalize(this.type)} - ${this.name}` } + protected getMachineId() { + return `${this.type}-${this.name}`; + } public notify = async (payload: WebhookPayload) => { this.emitEvent('notify', payload); @@ -295,6 +303,18 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i } } + protected async doParseCache(): Promise { + const cachedQueue = (await this.cache.cacheScrobble.get(`${this.getMachineId()}-queue`) as QueuedScrobble[] ?? []); + const cachedQLength = cachedQueue.length; + this.queuedScrobbles = cachedQueue; + + const cachedDead = (await this.cache.cacheScrobble.get(`${this.getMachineId()}-dead`) as DeadLetterScrobble[] ?? []); + const cachedDLength = cachedDead.length; + this.deadLetterScrobbles = cachedDead; + + return `Scrobbled from Cache: ${cachedQLength} Queue | ${cachedDLength} Dead Letter`; + } + protected async postInitialize(): Promise { const { options: { @@ -887,10 +907,12 @@ ${closestMatch.breakdowns.join('\n')}`, {leaf: ['Dupe Check']}); } this.logger.info(`Removed scrobble ${buildTrackString(this.deadLetterScrobbles[index].play)} from queue`, {leaf: 'Dead Letter'}); this.deadLetterScrobbles.splice(index, 1); + this.cache.cacheScrobble.set(`${this.getMachineId()}-dead`, this.deadLetterScrobbles); } removeDeadLetterScrobbles = () => { this.deadLetterScrobbles = []; + this.cache.cacheScrobble.set(`${this.getMachineId()}-dead`, []); this.logger.info('Removed all scrobbles from queue', {leaf: 'Dead Letter'}); } @@ -910,6 +932,7 @@ ${closestMatch.breakdowns.join('\n')}`, {leaf: ['Dupe Check']}); this.queuedScrobbles.push(queuedPlay); } this.queuedScrobbles.sort((a, b) => sortByOldestPlayDate(a.play, b.play)); + this.cache.cacheScrobble.set(`${this.getMachineId()}-queue`, this.queuedScrobbles); } protected addDeadLetterScrobble = (data: QueuedScrobble, error: (Error | string) = 'Unspecified error') => { @@ -923,6 +946,7 @@ ${closestMatch.breakdowns.join('\n')}`, {leaf: ['Dupe Check']}); this.deadLetterScrobbles.push(deadData); this.deadLetterScrobbles.sort((a, b) => sortByOldestPlayDate(a.play, b.play)); this.emitEvent('deadLetter', {dead: deadData}); + this.cache.cacheScrobble.set(`${this.getMachineId()}-dead`, this.deadLetterScrobbles); } queuePlayingNow = (data: PlayObject, source: SourceIdentifier) => { diff --git a/src/backend/tests/scrobbler/scrobblers.test.ts b/src/backend/tests/scrobbler/scrobblers.test.ts index be5c4a95..4755595f 100644 --- a/src/backend/tests/scrobbler/scrobblers.test.ts +++ b/src/backend/tests/scrobbler/scrobblers.test.ts @@ -16,6 +16,8 @@ import MockDate from 'mockdate'; import { NowPlayingScrobbler, TestAuthScrobbler, TestScrobbler } from "./TestScrobbler.js"; import { PlayPlatformId } from '../../common/infrastructure/Atomic.js'; +import { getRoot } from '../../ioc.js'; +import { loggerTest } from '@foxxmd/logging'; chai.use(asPromised); @@ -29,6 +31,8 @@ const normalizedWithMixedDur = normalizePlays(mixedDurPlays, {initialDate: first const normalizedWithMixedDurOlder = normalizePlays(mixedDurPlays, {initialDate: olderFirstPlayDate}); +getRoot({cache: {scrobble: {provider: 'memory'}}, logger: loggerTest}); + const generateTestScrobbler = () => { const testScrobbler = new TestScrobbler(); testScrobbler.verboseOptions = {