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.
59 kB · 1270 lines
TypeScript
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271import { 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}
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(); 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 = e; this.logger.error(e); }) } return new Promise((resolve, reject) => 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; }const noopTransform = async (x) => x; 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(); 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 { lastReadyAt: undefined, 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 (([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[]> => { 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; }
getRecentPlays = async (hydrate: boolean = true): Promise<PlayObject[]> => { 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; }
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) { 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))}`, priority: 'error'}); } return; } } if(!(await this.onPollPreAuthCheck())) { return; } if(!(await this.onPollPostAuthCheck())) { return; }
this.setStatus('Starting polling...');
this.abortController = new AbortController(); this.pollingPromise = spawn(this.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', this.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 { options: { maxPollRetries = 5, retryMultiplier = DEFAULT_RETRY_MULTIPLIER, } } = this.config;
// 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.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.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('maxInterval' in this.config.data) { 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 => 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) if(playObjs[0].data.playDate.isAfter(this.lastActivityAt)) { this.lastActivityAt = playObjs[0].data.playDate; } 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.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'); this.ingressQueueAbortController = new AbortController(); this.ingressQueuePromise = spawn(this.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', this.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); queueState.error = undefined; 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; 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: { type: 'lt', date: playEntity.seenAt } }); 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; } } } 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 { await this.playRepo.updateById(existing.id, {updatedAt: dayjs()}); playEntity.state = 'duped'; events.push(stateChangeToPlayEvent({state: 'duped'})); playEntity.parentId = existing.id; }
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'})); events.push(queueCompletionStateToPlayEvent({...queueState, queueStatus: QUEUE_STATUS_FAILED, error: generateLoggableAbortReason('Interrupted by abort signal', this.ingressQueueAbortController.signal)})); 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})); } throw new PlayProcessingError(e, {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;
this.deadQueueAbortController = new AbortController(); this.deadQueuePromise = spawn(this.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', this.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('interval' in this.config.data) { interval = this.config.data.interval; } return interval; }
protected getMaxBackoff() { let maxInterval = DEFAULT_POLLING_MAX_INTERVAL;
if('maxInterval' in this.config.data) { 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'); const cLogger = await componentFileLogger(this.type, this.name, true, logConfig); this.componentLogger = childLogger(cLogger, this.logger.labels); stream.on('data', (d: LogDataPretty) => { const {level, msg, line, labels, ...rest} = d; if(d.labels.includes(this.loggerLabel)) { this.componentLogger[this.componentLogger.levels.labels[d.level]]({...rest, labels: difference(labels, this.logger.labels)}, msg); } }); } }}