diff --git a/src/backend/index.ts b/src/backend/index.ts index 0e2dc7f0..d5e9615a 100644 --- a/src/backend/index.ts +++ b/src/backend/index.ts @@ -16,8 +16,6 @@ import { appLogger, initLogger as getInitLogger } from "./common/logging.js"; import { getRoot } from "./ioc.js"; import { parseVersion } from "./version.js"; import { initServer } from "./server/index.js"; -import { createHeartbeatClientsTask } from "./tasks/heartbeatClients.js"; -import { createHeartbeatSourcesTask } from "./tasks/heartbeatSources.js"; import { isDebugMode, parseBool, retry, sleep } from "./utils.js"; import { readJson } from './utils/DataUtils.js'; import ScrobbleClients from './scrobblers/ScrobbleClients.js'; @@ -157,7 +155,7 @@ const configDir = process.env.CONFIG_DIR || path.resolve(projectDir, `./config`) } for(const c of scrobbleClients.clients) { - c.initHeartbeat(); + c.initTasks(); const res = await Promise.race([ sleep(2200), (async () => { @@ -168,16 +166,25 @@ const configDir = process.env.CONFIG_DIR || path.resolve(projectDir, `./config`) })() ]); if(res === undefined) { - logger.debug(`Not waiting for ${c.name} to finish init, moving on to the next client...`); + logger.debug(`Not waiting for Client ${c.name} to finish init, moving on to the next Client...`); } } - const sourceTask = createHeartbeatSourcesTask(scrobbleSources, logger); - scheduler.addSimpleIntervalJob(new SimpleIntervalJob({ - minutes: 20, - runImmediately: true - }, sourceTask, {id: 'sources_heart'})); - logger.debug('Added Source Heartbeat task to scheduler'); + for(const c of scrobbleSources.sources) { + c.initTasks(); + const res = await Promise.race([ + sleep(2200), + (async () => { + while(!c.isReady()) { + await sleep(400) + } + return true; + })() + ]); + if(res === undefined) { + logger.debug(`Not waiting for Source ${c.name} to finish init, moving on to the next Source...`); + } + } let runRetentionNow = parseBool(process.env.RETENTION_IMMEDIATE, false); diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index 99d5a26a..e69ec299 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -235,7 +235,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i await this.tryStopScrobbling(); } - public initHeartbeat(opts: {deadDelay?: number} = {}) { + 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({ diff --git a/src/backend/sources/AbstractSource.ts b/src/backend/sources/AbstractSource.ts index 54c03ca4..0c65f0c8 100644 --- a/src/backend/sources/AbstractSource.ts +++ b/src/backend/sources/AbstractSource.ts @@ -48,6 +48,7 @@ import { spawn, catchAbortError, isAbortError, rethrowAbortError, delay, forever import { AbortedError, generateLoggableAbortReason } from '../common/errors/MSErrors.js'; import { DrizzlePlayRepository, playToRepositoryCreatePlayOpts, queryArgsFromRequest, QueryPlaysOpts, RequestPlayQuery } from '../common/database/drizzle/repositories/PlayRepository.js'; import { asPlay } from '../../core/PlayMarshalUtils.js'; +import { AsyncTask, SimpleIntervalJob, ToadScheduler } from 'toad-scheduler'; export interface RecentlyPlayedOptions { limit?: number @@ -90,6 +91,7 @@ export default abstract class AbstractSource extends AbstractComponent implement manualListening?: boolean + scheduler: ToadScheduler = new ToadScheduler(); emitter: EventEmitter; protected SCROBBLE_BACKLOG_COUNT: number = 30; @@ -146,6 +148,59 @@ export default abstract class AbstractSource extends AbstractComponent implement } } + public initTasks() { + 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 => { + return this.heartbeatTask().then(() => null).catch((err) => { + this.logger.error(err); + }); + }, + (err: Error) => { + this.logger.error(err); + } + ), {id: 'heartbeat'})); + } else { + this.logger.warn('Heartbeat task is already added to scheduler.'); + } + } + + protected async heartbeatTask(): Promise { + 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 { + await this.tryInitialize({force: false, notify: true, notifyTitle: 'Could not initialize automatically'}); + } catch (e) { + this.logger.error(new Error('Could not initialize automatically', {cause: e})); + return false; + } + + if('discoverDevices' in this && typeof this.discoverDevices === 'function') { + this.discoverDevices(); + } + + 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.poll({force: false, notify: true}).catch(e => this.logger.error(e)); + } + return true; + } + } + return true; + } + protected async postCache(): Promise { await super.postCache(); this.generateStaggerMappers(); diff --git a/src/backend/tasks/heartbeatClients.ts b/src/backend/tasks/heartbeatClients.ts deleted file mode 100644 index b8e02cf1..00000000 --- a/src/backend/tasks/heartbeatClients.ts +++ /dev/null @@ -1,59 +0,0 @@ -import { childLogger, Logger } from '@foxxmd/logging'; -import { PromisePool } from "@supercharge/promise-pool"; -import { AsyncTask } from "toad-scheduler"; -import ScrobbleClients from "../scrobblers/ScrobbleClients.js"; - -export const createHeartbeatClientsTask = (clients: ScrobbleClients, parentLogger: Logger) => { - const logger = childLogger(parentLogger, ['Heartbeat', 'Clients']); - - return new AsyncTask( - 'Heartbeat', - (): Promise => { - logger.verbose('Starting check...'); - return PromisePool - .withConcurrency(1) - .for(clients.clients) - .process(async (client) => { - if(!client.isReady()) { - if(!client.canAuthUnattended()) { - client.logger.warn({labels: 'Heartbeat'}, 'Client is not ready but will not try to initialize because auth state is not good and cannot be correct unattended.') - return 0; - } - try { - await client.tryInitialize({force: false, notify: true, notifyTitle: 'Could not initialize automatically'}); - } catch (e) { - client.logger.error(new Error('Could not initialize automatically', {cause: e})); - return 1; - } - } - - if(!client.canAuthUnattended()) { - client.logger.warn({label: 'Heartbeat'}, 'Should be monitoring scrobbles but will not attempt to start because auth state is not good and cannot be correct unattended.'); - return 0; - } - - await client.processDeadLetterQueue(); - if(!client.scrobbling) { - client.logger.info({labels: 'Heartbeat'}, 'Should be processing scrobbles! Attempting to restart scrobbling...'); - client.initScrobbleMonitoring(); - return 1; - } - }).then(({results, errors}) => { - logger.verbose(`Checked Dead letter queue for ${clients.clients.length} clients.`); - const restarted = results.reduce((acc, curr) => acc += curr, 0); - if (restarted > 0) { - logger.info(`Attempted to start ${restarted} clients that were not processing scrobbles.`); - } - if (errors.length > 0) { - logger.error(`Encountered errors!`); - for (const err of errors) { - logger.error(err); - } - } - }); - }, - (err: Error) => { - logger.error(err); - } - ); -} diff --git a/src/backend/tasks/heartbeatSources.ts b/src/backend/tasks/heartbeatSources.ts deleted file mode 100644 index 6317e4b3..00000000 --- a/src/backend/tasks/heartbeatSources.ts +++ /dev/null @@ -1,65 +0,0 @@ -import { childLogger, Logger } from '@foxxmd/logging'; -import { PromisePool } from "@supercharge/promise-pool"; -import { AsyncTask } from "toad-scheduler"; -import { ChromecastSource } from "../sources/ChromecastSource.js"; -import ScrobbleSources from "../sources/ScrobbleSources.js"; - -export const createHeartbeatSourcesTask = (sources: ScrobbleSources, parentLogger: Logger) => { - const logger = childLogger(parentLogger, ['Heartbeat', 'Sources']); - - return new AsyncTask( - 'Heartbeat', - (): Promise => { - logger.verbose('Starting check...'); - return PromisePool - .withConcurrency(1) - .for(sources.sources) - .process(async (source) => { - if(!source.isReady()) { - if(!source.canAuthUnattended()) { - source.logger.warn({label: 'Heartbeat'}, 'Source is not ready but will not try to initialize because auth state is not good and cannot be correct unattended.'); - return 0; - } - try { - await source.tryInitialize({force: false, notify: true, notifyTitle: 'Could not initialize automatically'}); - } catch (e) { - source.logger.error(new Error('Could not initialize source automatically', {cause: e})); - return 1; - } - } - - if(source.type === 'chromecast') { - (source as ChromecastSource).discoverDevices(); - } - - if (source.canPoll && !source.polling) { - if(!source.canAuthUnattended()) { - source.logger.warn({label: 'Heartbeat'}, 'Should be polling but will not attempt to start because auth state is not good and cannot be correct unattended.'); - return 0; - } else { - source.logger.info({label: 'Heartbeat'}, 'Should be polling, attempting to start polling...'); - source.poll({force: false, notify: true}).catch(e => source.logger.error(e)); - } - return 1; - } - - return 0; - }).then(({results, errors}) => { - logger.verbose(`Checked ${sources.sources.length} sources for start signals.`); - const restarted = results.reduce((acc, curr) => acc += curr, 0); - if (restarted > 0) { - logger.info(`Attempted to start ${restarted} sources.`); - } - if (errors.length > 0) { - logger.error(`Encountered errors!`); - for (const err of errors) { - logger.error(err); - } - } - }); - }, - (err: Error) => { - logger.error(err); - } - ); -} diff --git a/src/backend/tests/scrobbler/scrobblers.test.ts b/src/backend/tests/scrobbler/scrobblers.test.ts index 6bdadbc1..eac04f81 100644 --- a/src/backend/tests/scrobbler/scrobblers.test.ts +++ b/src/backend/tests/scrobbler/scrobblers.test.ts @@ -1040,7 +1040,7 @@ describe('Now Playing', function() { await using npScrobbler = new NowPlayingScrobbler(); npScrobbler.nowPlayingTaskInterval = 10; await npScrobbler.initialize(); - npScrobbler.initHeartbeat(); + npScrobbler.initTasks(); await npScrobbler.queuePlayingNow(generateSourcePlayerObj({play:generatePlay({}, {deviceId: genGroupIdStr(generatePlayPlatformId())})}), {type: 'jellyfin', name: 'test'}); @@ -1054,7 +1054,7 @@ describe('Now Playing', function() { await using npScrobbler = new NowPlayingScrobbler(); npScrobbler.nowPlayingTaskInterval = 10; await npScrobbler.initialize(); - npScrobbler.initHeartbeat(); + npScrobbler.initTasks(); const now = dayjs();