Something went wrong. Try again.
[READ-ONLY] Mirror of https://github.com/FoxxMD/multi-scrobbler. Scrobble plays from multiple sources to multiple clients docs.multi-scrobbler.app
deezer docker jellyfin koito lastfm listenbrainz maloja mopidy mpris music music-assistant plex scrobble self-hosted spotify subsonic tautulli youtube-music
Something went wrong. Try again.
61 kB · 1293 lines
TypeScript
at strict
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294import { childLogger, type LogDataPretty, type LogLevel } from '@foxxmd/logging';import dayjs, { type Dayjs } from "dayjs";import type { EventEmitter } from "events";import type { FixedSizeList } from "fixed-size-list";import { DEAD_LETTER_RETRIES_DEFAULT, DEAD_QUEUE, INGRESS_QUEUE, PARSED_FROM, type PlayMatchResult, type PlayObject, QUEUE_STATUS_COMPLETED, QUEUE_STATUS_FAILED, SOURCE_SOT } from "../../core/Atomic.ts";import { buildTrackString, capitalize, truncateStringToLength } from "../../core/StringUtils.ts";import AbstractComponent from "../common/AbstractComponent.ts";import { type Authenticatable, DEFAULT_POLLING_INTERVAL, DEFAULT_POLLING_MAX_INTERVAL, DEFAULT_RETRY_MULTIPLIER, type GroupedFixedPlays, type InternalConfig, type ProgressAwarePlayObject,} from "../common/infrastructure/Atomic.ts";import type {PARSED_FROM_TYPE, PlayUserId, QueueContext} from '../../core/Atomic.ts';import type {DeviceId} from '../../core/Atomic.ts';import type {SourceConfig} from '../common/infrastructure/config/source/sources.ts';import type {SourceType} from "../../core/Atomic.ts";import { TRANSFORM_HOOK } from "../../core/Transform.ts";import TupleMap from "../common/TupleMap.ts";import { difference, isDebugMode, pollingBackoff, sleep, sortByOldestPlayDate,} from "../utils.ts";import { sortByNewestPlayDate } from '../../core/PlayUtils.ts';import { formatNumber } from '../../core/DataUtils.ts';import { timeToHumanTimestamp } from "../../core/TimeUtils.ts";import { todayAwareFormat } from "../../core/TimeUtils.ts";import { getRoot } from '../ioc.ts';import { componentFileLogger } from '../common/logging.ts';import { messageWithCausesTruncatedDefault } from "../../core/ErrorUtils.ts";import { existingScrobble, type ExistingScrobbleOpts } from '../utils/PlayComparisonUtils.ts';import { consumeQueue } from '../utils/AsyncUtils.ts';import pMap from 'p-map';import type { Counter } from 'prom-client';import { spawn, isAbortError, delay, throwIfAborted, waitForEvent } from 'abort-controller-x';import { generateLoggableAbortReason, StageChangeError } from '../common/errors/MSErrors.ts';import { type QueryPlaysOpts, type RequestPlayQuery, type WithPlayRelation } from '../common/database/drizzle/repositories/PlayRepository.ts';import { asPlay } from '../../core/PlayMarshalUtils.ts';import { AsyncTask, SimpleIntervalJob, ToadScheduler } from 'toad-scheduler';import { COMPONENT_STATE, type ComponentSourceApiJson, type ComponentState, type PlayApiCommonDetailed } from '../../core/Api.ts';import type {PaginatedResponse} from "../../core/Api.ts";import type { PlaySelectWithQueueStates, PlayWith } from '../common/database/drizzle/drizzleTypes.ts';import { PLAY_EVENT_TYPE, type PlayEvent } from '../../core/PlayEvent.ts';import { dupeCheckToPlayEvent, queueCompletionStateToPlayEvent, stateChangeToPlayEvent, transformToPlayEvent } from '../common/database/drizzle/entityUtils.ts';import type { PlayProcessingResult } from '../common/infrastructure/PlayProcessing.ts';import { PlayProcessingError } from '../common/errors/PlayProcessingError.ts';
export interface RecentlyPlayedOptions { limit?: number formatted?: boolean
display?: boolean}
/** list is only guaranteed when hydrating, otherwise it is undefined if not already cached */type RecentPlaysGetter = { (hydrate?: true): Promise<PlayObject[]> (hydrate: boolean): Promise<PlayObject[] | undefined>}
export default abstract class AbstractSource extends AbstractComponent implements Authenticatable {
declare type: SourceType;
declare config: SourceConfig; clients: string[]; instantiatedAt: Dayjs; lastActivityAt: Dayjs;
multiPlatform: boolean = false;
localUrl: URL;
configDir: string;
canPoll: boolean = false; polling: boolean = false; canBacklog: boolean = false; protected deadQueueAbortController: AbortController | undefined; protected deadQueuePromise: Promise<void> | undefined; protected abortController: AbortController | undefined; protected pollingPromise: Promise<void> | undefined; stopPollingWaitInterval: number = 200; pollRetries: number = 0; tracksDiscovered: number = 0; tracksDiscoveredTotal: number = 0; deadLetterLength: number = 0; deadLetterQueued: number = 0;
queueIdleMs: number = 1000; queueConcurrency: number = 3;
protected isSleeping: boolean = false; protected wakeAt: Dayjs = dayjs();
supportsUpstreamRecentlyPlayed: boolean = false; supportsUpstreamNowPlaying: boolean = false; supportsManualListening: boolean = false;
scheduler: ToadScheduler = new ToadScheduler(); protected initDeadTimeout: NodeJS.Timeout | undefined;
protected SCROBBLE_BACKLOG_COUNT: number = 30;
protected recentDiscoveredPlays: GroupedFixedPlays = new TupleMap<DeviceId, PlayUserId, FixedSizeList<ProgressAwarePlayObject>>();
protected loggerLabel: string;
protected discoveredCounter: Counter;
declare protected componentType: 'source';
protected existingPlayOpts!: ExistingScrobbleOpts;
constructor(type: SourceType, name: string, config: SourceConfig, internal: InternalConfig, emitter: EventEmitter) { super(config); this.componentType = 'source'; const {clients = [] } = config; this.type = type; this.name = name; this.logger = childLogger(internal.logger, this.getIdentifier()); this.loggerLabel = this.getIdentifier(); this.dupeLogger = childLogger(this.logger, 'Dupe'); this.deadLogger = childLogger(this.logger, DEAD_QUEUE); this.config = config; this.clients = clients; this.logger.debug(`Scrobble To: ${this.clients.length === 0 ? 'All' : this.clients.join(' | ')}`); this.instantiatedAt = dayjs(); this.lastActivityAt = this.instantiatedAt; this.localUrl = internal.localUrl; this.configDir = internal.configDir; this.emitter = emitter; const metrics = getRoot().items.sourceMetics; this.discoveredCounter = metrics.discovered; this.queuedGauge = metrics.queued; this.deadLetterGauge = metrics.deadLetter;
this.existingPlayOpts = { logger: this.logger, transformRules: this.transformRules, transformPlay: this.transformPlay, existingSubmitted: async (_) => [undefined, undefined] } //this.existingPlay = (playObjPre: PlayObject, existingScrobbles: PlayObject[], log?: boolean) => existingScrobble(playObjPre, existingScrobbles, existingScrobbleOpts, log); }
existingPlay(playObjPre: PlayObject, existingScrobbles: PlayObject[], log?: boolean): Promise<PlayMatchResult> { return existingScrobble(playObjPre, existingScrobbles, this.existingPlayOpts, log); }
[Symbol.dispose]() { this.scheduler.stop(); for(const job of this.scheduler.getAllJobs()) { job.stop(); if(job.id !== undefined) { this.scheduler.removeById(job.id); } } } async [Symbol.asyncDispose]() { try { await this.stop({ reason: 'Instance is being destroyed' }); } catch (e) { this.logger.warn(e); } }
public initTasks(opts: {deadDelay?: number} = {}) { if(this.scheduler.existsById('heartbeat') === false) { this.logger.info('Adding Heartbeat Task and running immediately'); this.scheduler.addSimpleIntervalJob(new SimpleIntervalJob({ minutes: 20, runImmediately: true }, new AsyncTask( 'Heartbeat', (): Promise<any> => { return this.heartbeatTask().then(() => null).catch((err) => { this.errors.push(err); this.logger.error(err); }); }, (err: Error) => { this.logger.error(err); this.errors.push(err); } ), {id: 'heartbeat'})); } else { this.logger.verbose('Heartbeat task is already added to scheduler, running immediately instead'); const j = this.scheduler.getById('heartbeat') as SimpleIntervalJob; j.start(); }
if(this.scheduler.existsById('dead') === false && this.initDeadTimeout === undefined) { const deadDelay = opts.deadDelay ?? 120; this.logger.verbose(`Delaying Dead Scrobbler Processing Task by ${deadDelay} seconds`); this.initDeadTimeout = setTimeout(() => { this.logger.info('Adding Dead Scrobbler Processing Task and running immediately'); this.initDeadTimeout = undefined; this.scheduler.addSimpleIntervalJob(new SimpleIntervalJob({ minutes: 20, runImmediately: true }, new AsyncTask( 'Dead', (): Promise<any> => { if(this.isReady()) { return this.processDeadLetterQueue(undefined, 'Reprocessing bulk dead Plays by system').then(() => null).catch((e) => { this.warnings.push(e); this.logger.error(e); }) } return Promise.resolve(); }, (err: Error) => { this.warnings.push(err); this.logger.error(err); } ), {id: 'dead'})); }, deadDelay * 1000);
} else { if(this.initDeadTimeout !== undefined) { this.logger.verbose('Dead scrobble task timeout is already set'); } else { this.logger.verbose('Dead scrobble task is already added to the scheduler'); } } }
protected async heartbeatTask(): Promise<boolean> { if(!this.isReady()) { if(!this.canAuthUnattended()) { this.logger.warn({labels: 'Heartbeat'}, 'Source is not ready but will not try to initialize because auth state is not good and cannot be corrected unattended.') return false; } try { this.setStatus('Attempting to initialize...'); await this.initialize({force: false, notify: true, notifyTitle: 'Could not initialize automatically'}); } catch (e) { this.logger.error(new Error('Could not initialize automatically', {cause: e})); this.setStatus('Could not initialize automatically'); return false; }
if('discoverDevices' in this && typeof this.discoverDevices === 'function') { this.discoverDevices(); } } if(this.isReady()) { if(this.ingressQueuePromise === undefined) { this.setStatus('Starting discovery queue...'); await this.startDiscoveryQueue(); } if (this.canPoll && !this.polling) { if(!this.canAuthUnattended()) { this.logger.warn({labels: 'Heartbeat'}, 'Should be polling but will not attempt to start because auth state is not good and cannot be correct unattended.'); return false; } else { this.logger.info({labels: 'Heartbeat'}, 'Should be polling, attempting to start polling...'); this.setStatus('Attempting to start polling...'); this.poll({force: false, notify: true}).catch(e => this.logger.error(e)); } return true; } else if(!this.canPoll) { this.updateDates({lastReadyAt: dayjs(), force: true}); } } return true; }
public async start(opts: {forceInit?: boolean} = {}) { try { if (opts.forceInit) { if (!this.canAuthUnattended()) { this.logger.warn({ labels: 'Heartbeat' }, 'Source is not ready but will not try to initialize because auth state is not good and cannot be corrected unattended.') return false; } try { this.setStatus('Attempting to initialize...'); await this.initialize({ force: true, notify: true, notifyTitle: 'Could not initialize automatically' }); } catch (e) { this.logger.error(new Error('Could not initialize automatically', { cause: e })); this.setStatus('Could not initialize automatically'); return false; }
if ('discoverDevices' in this && typeof this.discoverDevices === 'function') { this.discoverDevices(); } } this.initTasks();
return true; } catch (e) { throw new StageChangeError('Failed to start', { cause: e }); } finally { this.emitComponentUpdate({state: this.getRunningState()}); } }
public async stop(opts: { reason?: string | Error } = {}) { try { if (this.canPoll) { await this.tryStopPolling(opts.reason); } await this.tryStopDiscoveryQueue(opts.reason); this.scheduler.stop(); for (const job of this.scheduler.getAllJobs()) { job.stop(); if(job.id !== undefined) { this.scheduler.removeById(job.id); } } this.setStatus('Stopped'); this.emitComponentUpdate<Partial<ComponentSourceApiJson>>({state: COMPONENT_STATE.STOPPED}); } catch (e) { this.emitComponentUpdate<Partial<ComponentSourceApiJson>>({state: this.getRunningState()}); throw new StageChangeError('Failed to stop', { cause: e }); } }
protected async postCache(): Promise<void> { await super.postCache(); }
protected async postDatabase(): Promise<void> { // this.playRepo = new DrizzlePlayRepository(this.db, {logger: this.logger}); // this.queueRepo = new DrizzleQueueRepository(this.db, {logger: this.logger}); // this.playEventsRepo = new DrizzlePlayEventsRepository(this.db, {logger: this.logger}); // this.playRepo.componentId = this.dbComponent.id; // this.queueRepo.componentId = this.dbComponent.id; const counts = await this.playRepo.getComponentPlayCountByState(); const discoveredCount = counts.find(x => x.state === 'discovered'); if(discoveredCount !== undefined) { this.tracksDiscoveredTotal = discoveredCount['count(*)']; } await this.updateQueueStats([INGRESS_QUEUE, DEAD_QUEUE]); }
public getRunningState(): ComponentState { if(this.scheduler.getAllJobs().length === 0) { return COMPONENT_STATE.STOPPED; }
const running = (this.canPoll && this.polling) || !this.canPoll;
if(running && !this.isMonitoring()) { return COMPONENT_STATE.IGNORED; } return running ? COMPONENT_STATE.RUNNING : COMPONENT_STATE.IDLE; }
protected getComponentApiData() { return { authType: this.authType, initialized: this.initializedOnce, hasAuth: this.requiresAuth, hasAuthInteraction: this.requiresAuthInteraction, authed: this.authed, } }
public getApiData(): ComponentSourceApiJson { return { lastImport: undefined, lastImportSuccess: undefined, ...super.getApiData(), ...this.getComponentApiData(), type: this.type, status: this.status, players: {}, tracksDiscovered: this.tracksDiscovered, deadLetterPlays: this.deadLetterQueued, queued: this.queuedLength, deadLetterPlaysTotal: this.deadLetterLength, sot: SOURCE_SOT.HISTORY, supportsUpstreamRecentlyPlayed: this.supportsUpstreamRecentlyPlayed, sleeping: this.getIsSleeping(), wakeAt: this.wakeAt !== undefined ? this.wakeAt.toISOString() : undefined, countLive: this.tracksDiscoveredTotal } }
getRecentlyPlayed = async (options: RecentlyPlayedOptions = {}): Promise<PlayObject[]> => []
getUpstreamRecentlyPlayed = async (options: RecentlyPlayedOptions = {}): Promise<PlayObject[]> => { throw new Error('Not implemented'); }
getUpstreamNowPlaying = async(): Promise<PlayObject[]> => { throw new Error('Not implemented'); }
// by default if the track was recently played it is valid // this is useful for sources where the track doesn't have complete information like Subsonic // TODO make this more descriptive? or move it elsewhere recentlyPlayedTrackIsValid = (playObj: PlayObject) => true
async findPreQueueExistingPlay(queueablePlay: PlayObject, context?: QueueContext & {isRetry?: boolean}) { /** * WHEN RUN IN A SOURCE COMPONENT * * for backlog and history plays we intentionally queue up plays that may have been processed before... * ...for backlog we do this to catch any plays that may have been missed during network outage or MS offline or just the source reporting new things * ...for history this is entirely how we "discover" new plays: we use MS's existing check logic to see find "new" plays on the same list (of history) that evolves over time * * for these two cases, for the majority of scenarios, we want to prune already processed plays from hitting the database * otherwise we are causing a lot of noise for duped plays IE history polls every minute and we don't want all 100+ already seen plays being persisted as duped every minute. * * ...for INGRESS we skip check because the assumption is whatever client is sending requests is very intentional and the user wants to see that their activity was recieved */ if (queueablePlay.meta.parsedFrom !== undefined && ([PARSED_FROM.history, PARSED_FROM.backlog] as PARSED_FROM_TYPE[]).includes(queueablePlay.meta.parsedFrom)) { // we should be adding Plays to the queue without any transforms // so run on "raw" play input // if we have seen a play with close temporality with the exact input hash then skip it entirely const cheapInputExisting = await this.playRepo.checkExisting(queueablePlay, { inputHash: queueablePlay, // existing should have been created *before* this play seenAt: { type: 'lt', date: dayjs(), } }); if (cheapInputExisting !== undefined) { if (isDebugMode()) { // log to trace for backlog for some visibility into what was pruned // this is fine noise-wise since this only happens when a component it (re)started // // for history we only want to do this if debugmode is enabled // TODO implement debugmode per component so global debug doesn't cause noise if this isn't the component that is being debugged this.logger.trace(`Not adding ${buildTrackString(queueablePlay)} to queue because it already exists in db as Play ${cheapInputExisting.uid}`); } return cheapInputExisting; } } return; }
getFlatRecentlyDiscoveredPlays = async (): Promise<PlayObject[]> => { const list: PlayObject[] = await this.getRecentlyDiscoveredPlays(); return list.sort(sortByNewestPlayDate); }
getRecentPlaysApi = async (query: RequestPlayQuery) => { const res = await this.playRepo.findPlays({ limit: 100, }); return res.map((x) => { const {id, ...rest} = x; return rest; }) }
protected recentDiscoveredCacheKey = () => { return `recentDiscovered-${this.dbComponent.id}`; }
protected recentCacheKey = () => { return `recent-${this.dbComponent.id}`; }
getRecentlyDiscoveredPlays = (async (hydrate: boolean = true): Promise<PlayObject[] | undefined> => { const cacheKey = this.recentDiscoveredCacheKey(); let list = await this.cache.cacheDb.get<PlayObject[]>(cacheKey); if(list === undefined && hydrate) { list = (await this.playRepo.findPlays({ state: ['discovered'], order: 'desc', sort: 'playedAt', limit: 200 })).map(x => ({...asPlay(x.play), id: x.id, uid: x.uid})) list.sort(sortByOldestPlayDate); await this.cache.cacheDb.set<PlayObject[]>(cacheKey, list, '2m'); } return list; }) as RecentPlaysGetter;
getRecentPlays = (async (hydrate: boolean = true): Promise<PlayObject[] | undefined> => { const cacheKey = this.recentCacheKey(); let list = await this.cache.cacheDb.get<PlayObject[]>(cacheKey); if(list === undefined && hydrate) { list = (await this.playRepo.findPlays({ stateNot: ['queued'], order: 'desc', sort: 'playedAt', limit: 200 })).map(x => ({...asPlay(x.play), id: x.id, uid: x.uid})) list.sort(sortByOldestPlayDate); await this.cache.cacheDb.set<PlayObject[]>(cacheKey, list, '2m'); } return list; }) as RecentPlaysGetter;
async existingDiscovered(play: PlayObject): Promise<PlayMatchResult> { let list: PlayObject[] = await this.getRecentPlays(true); if(play.id !== undefined) { // don't want to return the same play (by id) when checking by source list = list.filter(x => x.id === undefined || (x.id !== play.id)); } return await this.existingPlay(play, list); // if(matchResults.match) { // return matchResults.closestMatchedPlay; // } // return undefined; }
protected scrobble = async (newDiscoveredPlays: PlayObject[], options: { forceRefresh?: boolean, [key: string]: any, discoverLocation?: 'backlog' | [key: string] } = {}) => {
if(newDiscoveredPlays.length > 0) { newDiscoveredPlays.sort(sortByOldestPlayDate); const postCompareMapped = await pMap(newDiscoveredPlays, async (x) => await this.transformPlay(x, TRANSFORM_HOOK.postCompare), {concurrency: 3}); const events: PlayEvent[] = []; for(const p of postCompareMapped) { const {lifecycle = []} = p; const psLifecycle = lifecycle.filter(x => x.hook === TRANSFORM_HOOK.postCompare); if(psLifecycle.length > 0 && p.id !== undefined) { events.push({...transformToPlayEvent(psLifecycle), playId: p.id, createdAt: dayjs()}); } } this.emitEvent('discoveredToScrobble', { data: postCompareMapped, options: { ...options, checkTime: newDiscoveredPlays[newDiscoveredPlays.length-1].data.playDate?.add(2, 'second'), scrobbleFrom: this.getIdentifier(), scrobbleTo: this.clients } }); this.setStatus(`Forwarded ${newDiscoveredPlays.length} new Plays to Clients${options.discoverLocation !== undefined ? ` from ${options.discoverLocation} ` : ''}`); } }
protected processBacklog = async (signal: AbortSignal) => { if (this.canBacklog) {
const { options: { scrobbleBacklog = true } = {} } = this.config;
if(scrobbleBacklog === false) { this.logger.info('Source is able to scrobble backlog but was it disabled by user.'); this.setStatus('Not scrobbling backlog because it was disabled by user'); return; }
this.logger.info('Discovering backlogged tracks from recently played API...'); this.setStatus('Discovering backlogged tracks from recently played API...'); let backlogPlays: PlayObject[]; const { scrobbleBacklogCount = this.SCROBBLE_BACKLOG_COUNT } = this.config.options || {}; let backlogLimit = scrobbleBacklogCount; if(backlogLimit > this.SCROBBLE_BACKLOG_COUNT) { this.logger.warn(`scrobbleBacklogCount (${scrobbleBacklogCount}) cannot be greater than max API limit (${this.SCROBBLE_BACKLOG_COUNT}), reverting to max...`); backlogLimit = this.SCROBBLE_BACKLOG_COUNT; } try { this.logger.verbose(`Fetching the last ${backlogLimit}${backlogLimit === this.SCROBBLE_BACKLOG_COUNT ? ' (max) ' : ''} listens to check for backlogging...`); backlogPlays = (await this.getBackloggedPlays({limit: backlogLimit})).map((x) => ({...x, meta: {...x.meta, parsedFrom: PARSED_FROM.backlog}})); signal.throwIfAborted(); } catch (e) { throw new Error('Error occurred while fetching backlogged plays', {cause: e}); } await this.queuePlay(backlogPlays); this.logger.info('Backlog Plays added to discovery queue.'); } return; }
protected getBackloggedPlays = async (options: RecentlyPlayedOptions): Promise<PlayObject[]> => { this.logger.debug('Backlogging not implemented'); return []; }
onPollPreAuthCheck = async (): Promise<boolean> => true
onPollPostAuthCheck = async (): Promise<boolean> => true
poll = async (options: {force?: boolean, notify?: boolean} = {}) => { const {force = false, notify = false} = options;
if(this.polling) { this.logger.error('Already polling!'); return; }
if(!this.isReady() || force) { try { await this.initialize(options); } catch (e) { const err = new Error('Cannot start polling because Source is not ready', {cause: e}); this.logger.error(err); this.setStatus('Polling Error'); this.replaceErrors(err, {predicate: (x) => x.message === err.message}); this.emitComponentUpdate<Partial<ComponentSourceApiJson>>({errors: this.errors}); if(notify) { await this.notify( {title: `Polling Error`, message: `Cannot start polling because Source is not ready: ${truncateStringToLength(500)(messageWithCausesTruncatedDefault(e as Error))}`, priority: 'error'}); } return; } } if(!(await this.onPollPreAuthCheck())) { return; } if(!(await this.onPollPostAuthCheck())) { return; }
this.setStatus('Starting polling...');
const abortController = new AbortController(); this.abortController = abortController; this.pollingPromise = spawn(abortController.signal, async (signal, { defer, fork }) => { defer(async () => { this.polling = false; this.isSleeping = false; this.emitEvent('statusChange', {status: 'Idle'}); this.emitComponentUpdate<Partial<ComponentSourceApiJson>>({state: COMPONENT_STATE.IDLE}); });
fork(async (fSignal) => { try { await this.processBacklog(fSignal); } catch (e) { throwIfAborted(fSignal); await this.notify({ title: `Polling Error`, message: 'Polling interrupted because error occurred while processing backlog.', priority: 'error' }); throw new Error('Polling interrupted because error occurred while processing backlog', { cause: e }); } }); await this.startPolling(signal); }).catch((e) => { const componentUpdate: Partial<ComponentSourceApiJson> = { state: COMPONENT_STATE.IDLE }; if (isAbortError(e)) { const err = generateLoggableAbortReason('Polling stopped', abortController.signal); this.logger.info(err); //this.logger.trace(e); componentUpdate.status = 'Polling cancelled'; } else { const err = new Error('Polling stopped with error', { cause: e }); this.logger.warn(err); componentUpdate.status = 'Polling stopped with error'; this.warnings.push(err); componentUpdate.warnings = this.warnings; } this.emitComponentUpdate<Partial<ComponentSourceApiJson>>(componentUpdate); }).finally(() => { this.abortController = undefined; this.pollingPromise = undefined; }); }
startPolling = async (signal: AbortSignal) => { signal.throwIfAborted(); // reset poll attempts if already previously run this.pollRetries = 0;
const { maxPollRetries = 5, retryMultiplier = DEFAULT_RETRY_MULTIPLIER, } = this.config.options ?? {};
// can't have negative retries! const maxRetries = Math.max(0, maxPollRetries);
if(this.polling === true) { this.logger.warn(`Already polling! Polling needs to be stopped before it can be started`); return; }
while (this.pollRetries <= maxRetries) { try { if(!this.isReady() && this.buildOK) { this.logger.verbose(`Source is no longer ready! Will attempt to reinitialize => Connection OK: ${this.connectionOK} | Auth OK: ${this.authed}`); const init = await this.initialize(); if(init === false) { throw new Error('Source failed reinitialization'); } signal.throwIfAborted(); } await this.doPolling(signal); } catch (e) { if(isAbortError(e)) { throw e; } if (this.pollRetries < maxRetries) { const delayFor = pollingBackoff(this.pollRetries + 1, retryMultiplier); this.logger.info(`Poll retries (${this.pollRetries}) less than max poll retries (${maxRetries}), restarting polling after ${delayFor} second delay...`); await this.notify({title: `Polling Retry`, message: `Encountered error while polling but retries (${this.pollRetries}) are less than max poll retries (${maxRetries}), restarting polling after ${delayFor} second delay. | Error: ${(e as Error).message}`, priority: 'warn'}); await sleep((delayFor) * 1000); this.pollRetries++; } else { this.logger.warn(`Poll retries (${this.pollRetries}) equal to max poll retries (${maxRetries}), stopping polling!`); await this.notify({title: `Polling Error`, message: `Encountered error while polling and retries (${this.pollRetries}) are equal to max poll retries (${maxRetries}), stopping polling!. | Error: ${(e as Error).message}`, priority: 'error'}); throw e; } } } }
tryStopPolling = async (reason?: string | Error) => { if(this.polling === false) { this.logger.warn(`Polling is already stopped!`); return true; } if(this.abortController === undefined) { this.logger.error('No abort controller found! Nothing to stop.'); return false; } this.abortController.abort(reason); let elapsed = 0; while(this.polling && elapsed < (10 * this.stopPollingWaitInterval)) { this.logger.verbose(`Waiting for polling stop signal to be acknowledged (waited ${formatNumber(elapsed/1000)}s)`); await sleep(this.stopPollingWaitInterval); elapsed += this.stopPollingWaitInterval; } if(this.polling) { this.logger.warn('Could not stop polling! Or polling signal was lost :('); return false; } return true; }
protected doPolling = async (signal: AbortSignal): Promise<true | undefined> => { signal.throwIfAborted();
this.logger.info('Polling started'); this.emitEvent('statusChange', {status: 'Running'}); this.emitComponentUpdate<Partial<ComponentSourceApiJson>>({state: COMPONENT_STATE.RUNNING}); await this.notify({title: `Polling Started`, message: 'Polling Started', priority: 'info'}); this.setStatus('Polling Started'); this.lastActivityAt = dayjs(); let checksOverThreshold = 0; const checkActiveFor = 120; let maxInterval = DEFAULT_POLLING_MAX_INTERVAL;
if(this.config.data !== undefined && 'maxInterval' in this.config.data && this.config.data.maxInterval !== undefined) { maxInterval = this.config.data.maxInterval; } let isInactive = false;
try { this.polling = true; while (true) { signal.throwIfAborted(); const pollFrom = dayjs(); let lastActivityLogLevel: LogLevel = 'trace';
let playObjs: PlayObject[]; try { playObjs = await this.getRecentlyPlayed({formatted: true}); } catch (e) { throw new Error('Error occurred while refreshing recently played', {cause: e}); } finally { signal.throwIfAborted(); }
const interval = this.getInterval(true); const maxBackoff = this.getMaxBackoff(); let sleepTime = interval;
if(playObjs.length > 0) { const now = dayjs().unix(); const closeToInterval = playObjs.some(x => x.data.playDate !== undefined && now - x.data.playDate.unix() < 5); if (playObjs.length > 0 && closeToInterval) { // because the interval check was so close to the play date we are going to delay client calls for a few secs // this way we don't accidentally scrobble ahead of any other clients (we always want to be behind so we can check for dups) // additionally -- it should be ok to have this in the for loop because played_at will only decrease (be further in the past) so we should only hit this once, hopefully
// make sure delay is less than possible polling interval const maxDelay = Math.min(10, interval * 0.75); this.logger.info(`Potential plays were discovered close to polling interval! Delaying scrobble clients refresh by ${maxDelay} seconds so other clients have time to scrobble first`); await sleep(maxDelay * 1000); } await this.queuePlay(playObjs); //newDiscovered = await this.discover(playObjs, {signal}); signal.throwIfAborted(); // this.scrobble(newDiscovered, // { // forceRefresh: closeToInterval // }); }
const activityMsgs: string[] = [];
if(playObjs.length > 0) { playObjs.sort(sortByNewestPlayDate); // only update date if the play date is after the current activity date (in the case of backlogged plays) const newestPlayDate = playObjs[0].data.playDate; if(newestPlayDate !== undefined && newestPlayDate.isAfter(this.lastActivityAt)) { this.lastActivityAt = newestPlayDate; } checksOverThreshold = 0; }
const activeThreshold = this.lastActivityAt.add(checkActiveFor, 's'); const inactiveFor = dayjs.duration(Math.abs(activeThreshold.diff(dayjs(), 'millisecond'))).humanize(false); const relativeActivity = dayjs.duration(this.lastActivityAt.diff(dayjs(), 'ms')); const humanRelativeActivity = relativeActivity.asSeconds() > -3 ? '' : ` (${timeToHumanTimestamp(relativeActivity)} ago)`; let friendlyInterval = `${formatNumber(sleepTime)}`; const friendlyLastFormat = todayAwareFormat(this.lastActivityAt); activityMsgs.push(`Last activity at ${friendlyLastFormat}${humanRelativeActivity}`); if (activeThreshold.isBefore(dayjs())) { friendlyInterval = formatNumber(maxInterval); checksOverThreshold++; if(sleepTime < maxInterval) { const checkVal = Math.min(checksOverThreshold, 1000); const backoff = Math.round(Math.max(Math.min(Math.min(checkVal, 1000) * 2 * (1.1 * checkVal), maxBackoff), 5)); friendlyInterval = `(${interval} + ${backoff})`; sleepTime = interval + backoff; } if(!isInactive) { lastActivityLogLevel = 'debug'; isInactive = true; } activityMsgs.push(`Inactive for ${inactiveFor} (last + ${checkActiveFor}s)`); } else if(isInactive) { activityMsgs.push('New Activity after inactive period'); lastActivityLogLevel = 'debug'; isInactive = false; } activityMsgs.push(`Next check in ${friendlyInterval}s`); this.logger[lastActivityLogLevel](activityMsgs.join(' | ')); this.setWakeAt(pollFrom.add(sleepTime, 'seconds')); this.setIsSleeping(true); this.emitComponentUpdate<Partial<ComponentSourceApiJson>>({sleeping: true, wakeAt: this.getWakeAt().toISOString()}) // set last active before we sleep this.updateDates({lastActiveAt: dayjs(), lastReadyAt: dayjs()}); while(dayjs().isBefore(this.getWakeAt())) { // check for polling status every half second and wait till wake up time await delay(signal, 500); } this.setIsSleeping(false); this.emitComponentUpdate<Partial<ComponentSourceApiJson>>({sleeping: false}); // if we have made it this far in the loop we can reset poll retries this.pollRetries = 0; } } catch (e) { if(!isAbortError(e)) { this.logger.error(new Error('Error occurred while polling', {cause: e})); } if((e as Error).message.includes('Status code: 401')) { this.authed = false; this.authFailure = true; } throw e; } finally { this.setIsSleeping(false); this.emitComponentUpdate<Partial<ComponentSourceApiJson>>({sleeping: false}); } }
startDiscoveryQueue = async () => { this.setStatus('Starting discovery queue processing'); const ingressQueueAbortController = new AbortController(); this.ingressQueueAbortController = ingressQueueAbortController; this.ingressQueuePromise = spawn(ingressQueueAbortController.signal, async (signal, { defer }) => { await this.processDiscoveryQueue(signal); }).catch((e) => { const componentUpdate: Partial<ComponentSourceApiJson> = { }; if (isAbortError(e)) { const err = generateLoggableAbortReason('Discovery queue processing stopped', ingressQueueAbortController.signal); this.logger.info(err); //this.logger.trace(e); componentUpdate.status = 'Discovery queue processing cancelled'; } else { const err = new Error('Scrobble processing stopped with error', { cause: e }); this.logger.warn(err); componentUpdate.status = 'Discovery queue stopped with error'; this.warnings.push(err); componentUpdate.warnings = this.warnings; } this.emitComponentUpdate<Partial<ComponentSourceApiJson>>(componentUpdate); }).finally(() => { this.ingressQueueAbortController = undefined; this.ingressQueuePromise = undefined; }); }
tryStopDiscoveryQueue = async (reason?: string | Error) => { if(this.ingressQueuePromise === undefined) { this.logger.verbose(`Discovery is already stopped`); return; } if(this.ingressQueueAbortController === undefined) { this.logger.error('No abort controller found! Nothing to stop.'); return false; } this.ingressQueueAbortController.abort(reason) let timePasssed = 0; while(this.ingressQueuePromise !== undefined && timePasssed < (this.stopPollingWaitInterval * 10)) { await sleep(this.stopPollingWaitInterval); timePasssed += this.stopPollingWaitInterval; this.logger.verbose(`Waiting for discovery processing stop signal to be acknowledged (waited ${timePasssed}ms)`); } if(this.ingressQueuePromise !== undefined) { throw new Error('Could not stop discovery processing! Or signal was lost'); } return true; }
protected processDiscoveryQueue = async (signal: AbortSignal) => { signal.throwIfAborted();
let taskFailures = 0;
const { options: { maxRequestRetries = 5, retryMultiplier = DEFAULT_RETRY_MULTIPLIER, } = {}, } = this.config; const maxRetries = Math.max(0, maxRequestRetries); const consumedIds = new Map<string, number>();
try { await consumeQueue( 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); this.logger.debug(`Delaying discovery of Play ${item.uid} task for ${delayFor}ms due to non-zero prior failures (${taskFailures})`); await sleep(delayFor, { signal }); } return this.handlePlayProcessing(item, signal) }, { concurrency: this.queueConcurrency, idleMs: this.queueIdleMs, signal, onSuccess: (item, queueId) => { consumedIds.delete(queueId); taskFailures = Math.max(taskFailures - 1, 0); }, onError: async (e: Error, queueId) => { consumedIds.delete(queueId); taskFailures++; this.emitter.emit('discoveryQueueError', e); if(taskFailures < maxRetries) { this.logger.info(`Discovery queue retries (${taskFailures}) less than max processing retries (${maxRetries}), continuing with processing...`); await this.notify({title: `Processing Retry`, message: `Encountered error while polling but retries (${taskFailures}) are less than max poll retries (${maxRetries}) so will continue. Error: ${e.message}`, priority: 'warn'}); } else { this.logger.warn(`Discovery queue retries (${taskFailures}) equal to max processing retries (${maxRetries}), stopping processing!`); await this.notify({title: `Processing Error`, message: `Encountered error while scrobble processing and retries (${taskFailures}) are equal to max processing retries (${maxRetries}), stopping processing!. | Error: ${e.message}`, priority: 'error'}); throw e; } }, onEmpty: () => { this.emitter.emit('queueEmptied'); } }, ); } catch (e) { throw e; }
}
//processQueueCurrentPlay async processPlay(playEntity: PlayWith<'queueStates' | 'events'>, signal?: AbortSignal): Promise<PlayProcessingResult> { signal?.throwIfAborted(); this.setStatus(`Processing Play ${playEntity.uid}`);
const queueState = playEntity.queueStates.find(x => x.queueName === INGRESS_QUEUE); if(queueState === undefined) { throw new Error(`Play ${playEntity.uid} does not have an ${INGRESS_QUEUE} queue state`); } queueState.error = null; const { context, } = queueState const { useCache = true, isRetry = false, transform = true, dupeCheck = true, } = context || {};
const isDead = queueState.retries > 0 || isRetry; const logger = isDead ? childLogger(this.logger, ['Dead', `Play ${playEntity.uid}`]) : childLogger(this.logger, [`Play ${playEntity.uid}`]); this.setStatus(`Processing ${isDead ? 'Dead ' : ''}Play ${playEntity.uid}`);
const events: Omit<PlayEvent, 'playId'>[] = []; try {
if(isRetry !== true && !playEntity.play.meta.wasMonitored) { logger.debug(`Not processing ${buildTrackString(playEntity.play)} because monitoring was disabled when Play was queued.`); playEntity.state = 'discarded'; events.push(stateChangeToPlayEvent({state: playEntity.state, reason: 'Not processing because monitoring was disabled when Play was queued'})); events.push(queueCompletionStateToPlayEvent({...queueState, queueStatus: QUEUE_STATUS_COMPLETED})); return {playEntity, queue: queueState, events}; } let preCompared = playEntity.play; if(transform) { const {lifecycle = [], ...rest} = await this.transformPlay(playEntity.play, TRANSFORM_HOOK.preCompare, {useCachedResult: useCache}); preCompared = rest; if(lifecycle.length > 0) { events.push({...transformToPlayEvent(lifecycle), createdAt: dayjs()}); } } let existing: PlayObject | undefined; if (dupeCheck) { // cheap check for existing const cheapExisting = await this.playRepo.checkExisting(preCompared, { notId: playEntity.id, // should only be a dupe if there are play entities that match that were created *before* this entity // that way we don't accidentally mark the "original" of some N number of duplicate inputs as a dupe as well // // IE the "oldest" play of a set of duplicates should not itself be marked as a dupe of the "newer" duplicates seenAt: playEntity.seenAt !== null ? { type: 'lt', date: playEntity.seenAt } : undefined }); if (cheapExisting !== undefined) { events.push(dupeCheckToPlayEvent({ match: true, reason: `Matched hash on existing Play ${cheapExisting.uid} with close temporality` })); existing = { ...cheapExisting.play, id: cheapExisting.id, uid: cheapExisting.uid }; } else { const matchRes = await this.existingDiscovered({...preCompared, id: playEntity.id, uid: playEntity.uid}); events.push(dupeCheckToPlayEvent({...matchRes, createdAt: dayjs().toISOString()})); if (matchRes.match) { existing = matchRes.closestMatchedPlay as PlayObject; } } } playEntity.play = preCompared; signal?.throwIfAborted(); if(existing === undefined) { playEntity.state = 'discovered'; //state = 'discovered'; events.push(stateChangeToPlayEvent({state: 'discovered'})); this.tracksDiscovered++; this.tracksDiscoveredTotal++ this.discoveredCounter.labels(this.getPrometheusLabels()).inc(); this.emitEvent('discovered', {play: preCompared}); await this.scrobble([{...playEntity.play, id: playEntity.id, uid: playEntity.uid}]); } else { // existing plays are always from the db so should always have an id if(existing.id !== undefined) { await this.playRepo.updateById(existing.id, {updatedAt: dayjs()}); } playEntity.state = 'duped'; events.push(stateChangeToPlayEvent({state: 'duped'})); playEntity.parentId = existing.id ?? null; }
const recentPlays = await this.getRecentPlays(false); // only need to update if its already in memory, // and better to update in-memory than clear cache so we aren't refetching from db on every discover if(recentPlays !== undefined) { recentPlays.push({...preCompared, id: playEntity.id, uid: playEntity.uid}); recentPlays.sort(sortByOldestPlayDate); this.cache.cacheDb.set(this.recentCacheKey(), recentPlays, '2m'); } if(playEntity.state === 'discovered') { const recentDiscoveredPlays = await this.getRecentlyDiscoveredPlays(false); if(recentDiscoveredPlays !== undefined) { recentDiscoveredPlays.push({...preCompared, id: playEntity.id, uid: playEntity.uid}); recentDiscoveredPlays.sort(sortByOldestPlayDate); this.cache.cacheDb.set(this.recentDiscoveredCacheKey(), recentDiscoveredPlays, '2m'); } } events.push(queueCompletionStateToPlayEvent({...queueState, queueStatus: QUEUE_STATUS_COMPLETED})); logger.info(`${capitalize(playEntity.state)} => ${buildTrackString(preCompared)}`); return {playEntity, events, queue: queueState}; } catch (e) { if(e instanceof PlayProcessingError) { throw e; } if(isAbortError(e)) { events.push(stateChangeToPlayEvent({state: 'failed'})); const abortSignal = signal ?? this.ingressQueueAbortController?.signal; events.push(queueCompletionStateToPlayEvent({...queueState, queueStatus: QUEUE_STATUS_FAILED, error: abortSignal !== undefined ? generateLoggableAbortReason('Interrupted by abort signal', abortSignal) : e as Error})); throw e; } if(!events.some(x => x.eventName === PLAY_EVENT_TYPE.playStateChange)) { events.push(stateChangeToPlayEvent({state: 'failed'})); playEntity.state = 'failed'; } if(!events.some(x => x.eventName === PLAY_EVENT_TYPE.queueStateChange)) { events.push(queueCompletionStateToPlayEvent({...queueState, queueStatus: QUEUE_STATUS_FAILED, error: e as Error})); } throw new PlayProcessingError(e as Error, {playEntity, queue: queueState, events, showStopping: true}); } }
processDeadLetterQueue = async (attemptWithRetries?: number, reason?: string, sync?: boolean) => {
const logger = childLogger(this.logger, ['Dead']); if (!(await this.isReady())) { logger.warn('Cannot process dead letter scrobbles because client is not ready.'); return; } if(this.deadQueueAbortController !== undefined) { logger.warn('Dead scrobbles are currently being processed, cannot restart right now.'); return; }
const { options: { deadLetterRetries = 3 } = {} } = this.config;
const retries = attemptWithRetries ?? deadLetterRetries;
const deadQueueAbortController = new AbortController(); this.deadQueueAbortController = deadQueueAbortController; this.deadQueuePromise = spawn(deadQueueAbortController.signal, async (signal, { defer, fork }) => {
//const processable = await this.queueRepo.getQueueCount(this.dbComponent.id, [INGRESS_QUEUE], [QUEUE_STATUS_FAILED], retries); const processableArgs: QueryPlaysOpts = {queues: [{queueName: INGRESS_QUEUE, queueStatus: QUEUE_STATUS_FAILED, retries}], with: ['queues']}; let processable = await this.playRepo.findPlaysPaginated(processableArgs); this.deadLetterQueued = processable.meta.total;
const total = await this.queueRepo.getQueueCount(this.dbComponent.id, [INGRESS_QUEUE], {queueStatus: [QUEUE_STATUS_FAILED], retries: 10000}); this.deadLetterLength = total; const queueStatus = `${processable.meta.total} of ${total} dead Plays have less than ${retries} retries, ${processable.meta.total === 0 ? 'will skip processing.': 'processing now...'}`; if (processable.meta.total === 0) { logger.verbose(queueStatus); return; } this.setStatus(`Queuing ${processable} Dead Plays...`); logger.info(queueStatus); let more = true; let offset = 0; while(more) { await this.queuePlay(processable.data, {reason}); more = processable.data.length === processable.meta.limit; if(more) { offset += processable.meta.limit; processable = await this.playRepo.findPlaysPaginated({...processableArgs, offset}); } } this.setStatus(`All processable Dead Plays have been queued`); logger.info(`All processable Dead Plays have been queued`);
if(sync) { await waitForEvent(signal,this.emitter,'queueEmptied'); this.setStatus(`Finished processing Dead Plays`); logger.info('Finished processing Dead Plays'); } }).catch((e) => { if (isAbortError(e)) { const err = generateLoggableAbortReason('Dead scrobble processing stopped', deadQueueAbortController.signal); this.logger.info(err); logger.trace(e) } else { logger.warn(new Error('Dead scrobble processing stopped with error', { cause: e })); } }).finally(() => { this.deadQueueAbortController = undefined; this.deadQueuePromise = undefined; }); }
protected setIsSleeping(sleeping: boolean) { this.isSleeping = sleeping; }
protected getIsSleeping() { return this.isSleeping; }
protected setWakeAt(dt: Dayjs) { this.wakeAt = dt; }
protected getWakeAt() { return this.wakeAt; }
protected getInterval(log?: boolean) { let interval = DEFAULT_POLLING_INTERVAL;
if(this.config.data !== undefined && 'interval' in this.config.data && this.config.data.interval !== undefined) { interval = this.config.data.interval; } return interval; }
protected getMaxBackoff() { let maxInterval = DEFAULT_POLLING_MAX_INTERVAL;
if(this.config.data !== undefined && 'maxInterval' in this.config.data && this.config.data.maxInterval !== undefined) { maxInterval = this.config.data.maxInterval; } return maxInterval - this.getInterval(); }
public async getPlaysPaginated(args: QueryPlaysOpts): Promise<PaginatedResponse<PlayApiCommonDetailed>> { const { limit, offset, with: withQuery = ['input','parent-input','queues'], ...rest } = args; const parsedLimit = limit !== undefined ? Number.parseInt(limit as unknown as string) : undefined; const parsedOffset = offset !== undefined ? Number.parseInt(offset as unknown as string) : undefined; return this.playRepo.findPlaysPaginated({limit: parsedLimit, offset: parsedOffset, with: withQuery, ...rest}); }
public async getPlaysPaginatedInternal(args: QueryPlaysOpts) { const { limit, offset, with: withQuery = ['input','parent-input','queues'], ...rest } = args; const parsedLimit = limit !== undefined ? Number.parseInt(limit as unknown as string) : undefined; const parsedOffset = offset !== undefined ? Number.parseInt(offset as unknown as string) : undefined; return this.playRepo.findPlaysPaginated({limit: parsedLimit, offset: parsedOffset, with: withQuery, ...rest}); }
public async getPlayApiResponse(uid: string, opts: {with?: WithPlayRelation[]} = {}): Promise<PlayApiCommonDetailed> { const { with: withQuery = ['input','parent-input','queues','events'], } = opts; return await this.playRepo.findByUid(uid, { with: withQuery as WithPlayRelation[] }) as unknown as PlayApiCommonDetailed; }
public async deletePlay(play: PlayWith<'children'>, children?: boolean): Promise<void> { if(children) { await this.playRepo.deleteByIds([play.id, ...(play.children ?? []).map(x => x.id)]); this.emitEvent('playDelete', {uid: play.uid}); for(const p of play.children) { this.emitEvent('playDelete', {uid: p.uid, componentId: p.componentId}); } } else { await this.playRepo.deleteById(play.id); this.emitEvent('playDelete', {uid: play.uid, componentId: play.componentId}); } }
public emitEvent = (eventName: string, payload: object = {}) => { this.emitter.emit(eventName, { type: this.type, name: this.name, componentId: this.dbComponent?.id, from: 'source', data: payload, }); }
public async destroy() { this.emitter.removeAllListeners(); }
protected async doBuildComponentLogger(): Promise<void> { if(this.config?.options?.logToFile) { this.logger.debug('Enabling component logger...'); const root = getRoot(); const stream = root.get('loggerStream'); const logConfig = root.get('loggingConfig'); if(stream === undefined) { this.logger.warn('No logger stream is available, cannot build component logger'); return; } const cLogger = await componentFileLogger(this.type, this.name, true, logConfig); const componentLogger = childLogger(cLogger, this.logger.labels); this.componentLogger = componentLogger; stream.on('data', (d: LogDataPretty) => { const {level, msg, line, labels, ...rest} = d; if(d.labels.includes(this.loggerLabel)) { componentLogger[componentLogger.levels.labels[d.level] as LogLevel]({...rest, labels: difference(labels, this.logger.labels)}, msg as string); } }); } }}