diff --git a/src/backend/common/AbstractComponent.ts b/src/backend/common/AbstractComponent.ts index 1a2bef36..acbdd822 100644 --- a/src/backend/common/AbstractComponent.ts +++ b/src/backend/common/AbstractComponent.ts @@ -3,7 +3,7 @@ import { cacheFunctions, } from "@foxxmd/regex-buddy-core"; import type EventEmitter from "events"; -import type {ComponentType, LifecycleInput, LifecycleStep, PlayData, PlayObject} from "../../core/Atomic.ts"; +import {MONITORING_ORIGIN_SYSTEM, MONITORING_ORIGIN_USER, type ComponentType, type LifecycleInput, type LifecycleStep, type PlayData, type PlayObject} from "../../core/Atomic.ts"; import { buildTrackString } from "../../core/StringUtils.ts"; import type {CommonClientConfig} from "./infrastructure/config/client/index.ts"; import type {CommonSourceConfig} from "./infrastructure/config/source/index.ts"; @@ -33,7 +33,7 @@ import { getRetentionCompactAfterFromEnv, getRetentionDeleteAfterFromEnv, isComp import type {DbConcrete} from "./database/drizzle/drizzleUtils.ts"; import type {ComponentSelect} from "./database/drizzle/drizzleTypes.ts"; import { DrizzlePlayRepository } from "./database/drizzle/repositories/PlayRepository.ts"; -import type {ClientType} from "../../core/Atomic.ts"; +import type {ClientType, MonitoringStatus} from "../../core/Atomic.ts"; import type {SourceType} from "../../core/Atomic.ts"; import { DrizzleComponentRepository } from "./database/drizzle/repositories/ComponentRepository.ts"; import dayjs from "dayjs"; @@ -60,6 +60,9 @@ export default abstract class AbstractComponent extends AbstractInitializable { status: string = 'Waiting to initialize...'; emitter: EventEmitter; + monitoringActivity?: boolean | undefined; + monitoringActivityDefault: boolean = true; + protected componentType: ComponentType; type: ClientType | SourceType; name: string; @@ -540,6 +543,7 @@ export default abstract class AbstractComponent extends AbstractInitializable { name: this.dbComponent.name, state, mode: this.dbComponent.mode, + monitoringStatus: this.getMonitoringStatus(), countNonLive: this.dbComponent.countNonLive, createdAt: this.dbComponent.createdAt?.toISOString(), lastReadyAt: this.dbComponent.lastReadyAt?.toISOString(), @@ -585,4 +589,15 @@ export default abstract class AbstractComponent extends AbstractInitializable { this.status = status; this.emitComponentUpdate({status}); } + + public getSystemMonitoring = (): boolean => this.config.options?.autoMonitor ?? this.getSystemDefaultMonitoring(); + + protected getSystemDefaultMonitoring = (): boolean => this.monitoringActivityDefault; + + public isMonitoring = (): boolean => this.monitoringActivity ?? this.getSystemMonitoring(); + + public getMonitoringStatus = (): MonitoringStatus => ({ + monitoring: this.isMonitoring(), + origin: this.monitoringActivity !== undefined ? MONITORING_ORIGIN_USER : MONITORING_ORIGIN_SYSTEM + }) } diff --git a/src/backend/common/infrastructure/config/client/index.ts b/src/backend/common/infrastructure/config/client/index.ts index b85b8ca5..d6df33ba 100644 --- a/src/backend/common/infrastructure/config/client/index.ts +++ b/src/backend/common/infrastructure/config/client/index.ts @@ -1,6 +1,6 @@ import type {DurationValue} from "../../Atomic.ts"; import type {PlayTransformOptions} from "../../../../../core/Transform.ts"; -import type {CommonConfig, RequestRetryOptions} from "../common.ts"; +import type {CommonConfig, MonitorOptions, RequestRetryOptions} from "../common.ts"; import type {RetentionConfig} from "../database.ts"; /** @@ -83,7 +83,7 @@ export interface NowPlayingOptions { nowPlaying?: boolean | string[] } -export interface CommonClientOptions extends RequestRetryOptions, UpstreamRefreshOptions { +export interface CommonClientOptions extends RequestRetryOptions, UpstreamRefreshOptions, MonitorOptions { /** * Check client for an existing scrobble at the same recorded time as the "new" track to be scrobbled. If an existing scrobble is found this track is not track scrobbled. diff --git a/src/backend/common/infrastructure/config/common.ts b/src/backend/common/infrastructure/config/common.ts index 5868e7d9..fe6af15f 100644 --- a/src/backend/common/infrastructure/config/common.ts +++ b/src/backend/common/infrastructure/config/common.ts @@ -88,3 +88,12 @@ export interface PollingOptions { orphanedAfter?: number } +export interface MonitorOptions { + /** + * Set the default behavior for wether this component should automatically monitor any activity, or scrobble, it encounters + * + * @default true + * @examples [true, false] + */ + autoMonitor?: boolean +} diff --git a/src/backend/common/infrastructure/config/source/azuracast.ts b/src/backend/common/infrastructure/config/source/azuracast.ts index 16738997..0abc45c1 100644 --- a/src/backend/common/infrastructure/config/source/azuracast.ts +++ b/src/backend/common/infrastructure/config/source/azuracast.ts @@ -1,4 +1,4 @@ -import type {CommonSourceConfig, CommonSourceData, CommonSourceOptions, ManualListeningOptions} from "./index.ts"; +import type {CommonSourceConfig, CommonSourceData, CommonSourceOptions} from "./index.ts"; export interface AzuraStationInfoResponse { id: string @@ -93,7 +93,7 @@ export interface AzuracastData extends CommonSourceData { apiKey?: string } -export interface AzuracastSourceoptions extends CommonSourceOptions, ManualListeningOptions { +export interface AzuracastSourceoptions extends CommonSourceOptions { } diff --git a/src/backend/common/infrastructure/config/source/icecast.ts b/src/backend/common/infrastructure/config/source/icecast.ts index daf3dcef..61198f7a 100644 --- a/src/backend/common/infrastructure/config/source/icecast.ts +++ b/src/backend/common/infrastructure/config/source/icecast.ts @@ -1,4 +1,4 @@ -import type {CommonSourceConfig, CommonSourceData, CommonSourceOptions, ManualListeningOptions} from "./index.ts"; +import type {CommonSourceConfig, CommonSourceData, CommonSourceOptions} from "./index.ts"; export interface IcecastMetadata { @@ -34,7 +34,7 @@ export interface IcecastData extends CommonSourceData, IcecastOptions { url: string } -export interface IcecastSourceOptions extends CommonSourceOptions, ManualListeningOptions { +export interface IcecastSourceOptions extends CommonSourceOptions { } export interface IcecastSourceConfig extends CommonSourceConfig { diff --git a/src/backend/common/infrastructure/config/source/index.ts b/src/backend/common/infrastructure/config/source/index.ts index b58ab78f..f290368f 100644 --- a/src/backend/common/infrastructure/config/source/index.ts +++ b/src/backend/common/infrastructure/config/source/index.ts @@ -1,7 +1,7 @@ import type { FileLogOptions, LogLevel } from "@foxxmd/logging"; import type { PlayTransformOptions } from "../../../../../core/Transform.ts"; -import type { CommonConfig, RequestRetryOptions } from "../common.ts"; +import type { CommonConfig, MonitorOptions, RequestRetryOptions } from "../common.ts"; import type { RetentionConfig } from "../database.ts"; import type { DurationValue } from "../../Atomic.ts"; @@ -44,7 +44,7 @@ export interface ScrobbleThresholds { percent?: number | null } -export interface CommonSourceOptions extends SourceRetryOptions { +export interface CommonSourceOptions extends SourceRetryOptions, MonitorOptions { /** * * If this source has INGRESS to MS (sends a payload, rather than MS GETTING requesting a payload) then setting this option to true will make MS log the payload JSON to DEBUG output * * If this source is POLLING then it will log the raw data for each unique track/response the first time it is seen @@ -114,15 +114,6 @@ export interface CommonSourceOptions extends SourceRetryOptions { retention?: RetentionConfig } -export interface ManualListeningOptions { - /** - * For Sources that support manual listening, should MS default to scrobbling when no user interaction has occurred? - * - * If not specified MS will use a Source's specific behavior, see Source's documentation. - */ - systemScrobble?: boolean -} - export interface CommonSourceData { } diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index af5c35a1..f06a41ae 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -1143,6 +1143,14 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i let deadQueueEntity: QueueStateSelect; try { + + // want to fail scrobbles that are being *ingested* by component + // but not those that have already failed but will be dead queue processed + // since dead queued scrobbles are likely from a different time period that had monitoring + const monitoringStatus = this.getMonitoringStatus(); + if(!monitoringStatus.monitoring) { + throw new SimpleError(`Monitoring is disabled by ${capitalize(monitoringStatus.origin)}`); + } if (this.upstreamRefresh.refreshEnabled) { try { diff --git a/src/backend/server/api.ts b/src/backend/server/api.ts index 12f97e16..921a2802 100644 --- a/src/backend/server/api.ts +++ b/src/backend/server/api.ts @@ -380,9 +380,9 @@ export const setupApi = (app: Express, logger: Logger, appLoggerStream: PassThro players: 'players' in x ? (x as MemorySource).playersToObject() as unknown as Record : {}, sot: ('playerSourceOfTruth' in x) ? x.playerSourceOfTruth as SOURCE_SOT_TYPES : SOURCE_SOT.HISTORY, supportsUpstreamRecentlyPlayed: x.supportsUpstreamRecentlyPlayed, - supportsManualListening: x.supportsManualListening, - manualListening: x.manualListening, - systemListeningBehavior: x.getSystemListeningBehavior(), + supportsManualListening: true, + manualListening: x.monitoringActivity, + systemListeningBehavior: x.getSystemMonitoring(), ...x.additionalApiData() }; if(!x.isReady()) { diff --git a/src/backend/sources/AbstractSource.ts b/src/backend/sources/AbstractSource.ts index 966f3bb8..683c3723 100644 --- a/src/backend/sources/AbstractSource.ts +++ b/src/backend/sources/AbstractSource.ts @@ -14,7 +14,7 @@ import { type InternalConfig, type ProgressAwarePlayObject, } from "../common/infrastructure/Atomic.ts"; -import type {PlayUserId} from '../../core/Atomic.ts'; +import type {PlayState, PlayUserId} 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"; @@ -46,6 +46,7 @@ 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 { PlayWith } from '../common/database/drizzle/drizzleTypes.ts'; export interface RecentlyPlayedOptions { limit?: number @@ -293,22 +294,12 @@ export default abstract class AbstractSource extends AbstractComponent implement tracksDiscovered: this.tracksDiscovered, sot: SOURCE_SOT.HISTORY, supportsUpstreamRecentlyPlayed: this.supportsUpstreamRecentlyPlayed, - supportsManualListening: this.supportsManualListening, - manualListening: this.manualListening, - systemListeningBehavior: this.getSystemListeningBehavior(), sleeping: this.getIsSleeping(), wakeAt: this.wakeAt !== undefined ? this.wakeAt.toISOString() : undefined, countLive: this.tracksDiscoveredTotal } } - getSystemListeningBehavior = (): boolean | undefined => { - if(this.supportsManualListening) { - return this.config.options !== undefined && 'systemScrobble' in this.config.options ? this.config.options?.systemScrobble : undefined; - } - return undefined; - } - getRecentlyPlayed = async (options: RecentlyPlayedOptions = {}): Promise => [] getUpstreamRecentlyPlayed = async (options: RecentlyPlayedOptions = {}): Promise => { @@ -324,8 +315,14 @@ export default abstract class AbstractSource extends AbstractComponent implement // TODO make this more descriptive? or move it elsewhere recentlyPlayedTrackIsValid = (playObj: PlayObject) => true - protected addPlayToDiscovered = async (play: PlayObject): Promise => { - const playRow = await this.playRepo.createPlays([(playToRepositoryCreatePlayOpts({play, state: 'discovered'}))]); + protected addPlayToDB = async (play: PlayObject): Promise> => { + const monitorStatus = this.getMonitoringStatus(); + let state: PlayState = 'discovered'; + if(!monitorStatus.monitoring) { + this.logger.debug(`Not adding ${buildTrackString(play)} as discovered because monitoring is disabled by ${capitalize(monitorStatus.origin)}`); + state = 'discarded'; + } + const playRow = await this.playRepo.createPlays([(playToRepositoryCreatePlayOpts({play, state}))]); const recentPlays = await this.getRecentlyDiscoveredPlays(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 @@ -334,15 +331,17 @@ export default abstract class AbstractSource extends AbstractComponent implement recentPlays.sort(sortByOldestPlayDate); this.cache.cacheDb.set(this.recentDiscoveredCacheKey(), recentPlays, '2m'); } - this.tracksDiscovered++; - this.tracksDiscoveredTotal++ - this.logger.info(`Discovered => ${buildTrackString(play)}`); + if(state === 'discovered') { + this.tracksDiscovered++; + this.tracksDiscoveredTotal++ + this.discoveredCounter.labels(this.getPrometheusLabels()).inc(); + } + this.logger.info(`${capitalize(state)} => ${buildTrackString(play)}`); this.emitEvent('discovered', {play}); this.emitPlayInsert({...playRow[0], queueStates: []} as unknown as PlayApiCommonDetailed); - this.discoveredCounter.labels(this.getPrometheusLabels()).inc(); - play.id = playRow[0].id; - play.uid = playRow[0].uid; - return play; + playRow[0].play.id = playRow[0].id; + playRow[0].play.uid = playRow[0].uid; + return playRow[0]; } getFlatRecentlyDiscoveredPlays = async (): Promise => { @@ -398,8 +397,10 @@ export default abstract class AbstractSource extends AbstractComponent implement const existing = await this.existingDiscovered(play); if(existing === undefined) { options.signal?.throwIfAborted() - const hydratedPlay = await this.addPlayToDiscovered(play); - newDiscoveredPlays.push(hydratedPlay); + const hydratedPlay = await this.addPlayToDB(play); + if(hydratedPlay.state === 'discovered') { + newDiscoveredPlays.push(hydratedPlay.play); + } } else { this.playRepo.updateById(existing.id, {updatedAt: dayjs()}); } @@ -419,26 +420,10 @@ export default abstract class AbstractSource extends AbstractComponent implement return newDiscoveredPlays; } - protected shouldScrobble = (discoverLocation?: 'backlog' | [key: string]) => { - if(this.supportsManualListening && discoverLocation !== 'backlog') { - const manualFlag = this.manualListening ?? this.getSystemListeningBehavior() ?? true; - if(manualFlag === false) { - this.logger.debug(`NOT scrobbling because Should Scrobble is FALSE (${this.manualListening === false ? 'user' : 'system'})`); - return false; - } - } - return true; - } - protected scrobble = async (newDiscoveredPlays: PlayObject[], options: { forceRefresh?: boolean, [key: string]: any, discoverLocation?: 'backlog' | [key: string] } = {}) => { if(newDiscoveredPlays.length > 0) { - if(!this.shouldScrobble(options.discoverLocation)) { - await this.playRepo.setStateById('discarded', newDiscoveredPlays.map(x => x.id)); - this.setStatus(`Discarded ${newDiscoveredPlays} new Plays${options.discoverLocation !== undefined ? ` from ${options.discoverLocation} ` : ''}`); - return; - } newDiscoveredPlays.sort(sortByOldestPlayDate); this.emitter.emit('discoveredToScrobble', { data: await pMap(newDiscoveredPlays, this.staggerMappers.postCompare(async (x) => await this.transformPlay(x, TRANSFORM_HOOK.postCompare)), {concurrency: 3}), @@ -470,7 +455,7 @@ export default abstract class AbstractSource extends AbstractComponent implement this.logger.info('Discovering backlogged tracks from recently played API...'); this.setStatus('Discovering backlogged tracks from recently played API...'); - let backlogPlays: PlayObject[] = []; + let backlogPlays: PlayObject[]; const { scrobbleBacklogCount = this.SCROBBLE_BACKLOG_COUNT } = this.config.options || {}; @@ -652,11 +637,8 @@ export default abstract class AbstractSource extends AbstractComponent implement } this.abortController.abort(reason); let elapsed = 0; - let lastlog: Dayjs; while(this.polling && elapsed < (10 * this.stopPollingWaitInterval)) { - if(lastlog === undefined || dayjs().diff(lastlog, 's') >= 2) { - this.logger.verbose(`Waiting for polling stop signal to be acknowledged (waited ${formatNumber(elapsed/1000)}s)`); - } + this.logger.verbose(`Waiting for polling stop signal to be acknowledged (waited ${formatNumber(elapsed/1000)}s)`); await sleep(this.stopPollingWaitInterval); elapsed += this.stopPollingWaitInterval; } diff --git a/src/backend/sources/AppleMusicSource.ts b/src/backend/sources/AppleMusicSource.ts index 36099ce3..a595423e 100644 --- a/src/backend/sources/AppleMusicSource.ts +++ b/src/backend/sources/AppleMusicSource.ts @@ -360,7 +360,7 @@ export default class AppleMusicSource extends AbstractSource { reversedPlays.reverse(); for(const refPlay of reversedPlays) { - await this.addPlayToDiscovered(refPlay); + await this.addPlayToDB(refPlay); } } return true; diff --git a/src/backend/sources/AzuracastSource.ts b/src/backend/sources/AzuracastSource.ts index b569f288..d4d3f393 100644 --- a/src/backend/sources/AzuracastSource.ts +++ b/src/backend/sources/AzuracastSource.ts @@ -28,25 +28,11 @@ export class AzuracastSource extends MemorySource { wsNowPlaying: AzuraStationResponse wsCurrenTime: number = 0; client!: WS; + override monitoringActivityDefault = false; constructor(name: any, config: AzuracastSourceConfig, internal: InternalConfig, emitter: EventEmitter) { - const { - data = {}, - options = {}, - } = config; - const { - ...rest - } = data; - - const { - data: { - monitorWhenListeners, - monitorWhenLive - } = {} - } = config; - - super('azuracast', name, { ...config, options: {systemScrobble: monitorWhenListeners !== undefined || monitorWhenLive === true, ...options}, data: { ...rest } }, internal, emitter); + super('azuracast', name, config, internal, emitter); this.requiresAuth = false; @@ -228,6 +214,14 @@ export class AzuracastSource extends MemorySource { return await this.processRecentPlays([playerState]); } + protected getSystemDefaultMonitoring = (): boolean => { + const { + monitorWhenLive, + monitorWhenListeners + } = this.config.data; + return monitorWhenLive !== undefined || monitorWhenListeners !== undefined; + } + } const formatPlayObj = (obj: AzuraNowPlayingResponse, options: FormatPlayObjectOptions = {}): PlayObject => { diff --git a/src/backend/sources/IcecastSource.ts b/src/backend/sources/IcecastSource.ts index 5027f440..2b8b134d 100644 --- a/src/backend/sources/IcecastSource.ts +++ b/src/backend/sources/IcecastSource.ts @@ -30,6 +30,8 @@ export class IcecastSource extends MemorySource { streamError?: Error; streaming: boolean = false; + override monitoringActivityDefault = false; + constructor(name: any, config: IcecastSourceConfig, internal: InternalConfig, emitter: EventEmitter) { const { data, @@ -38,7 +40,7 @@ export class IcecastSource extends MemorySource { const { ...rest } = data || {}; - super('icecast', name, { ...config, options: {systemScrobble: false, ...options}, data: { ...rest } }, internal, emitter); + super('icecast', name, config, internal, emitter); this.requiresAuth = false; this.canPoll = true; diff --git a/src/backend/sources/YTMusicSource.ts b/src/backend/sources/YTMusicSource.ts index 6468c980..ed8b0ce8 100644 --- a/src/backend/sources/YTMusicSource.ts +++ b/src/backend/sources/YTMusicSource.ts @@ -691,7 +691,7 @@ ${humanDiff}`; // and add to discovered since its empty for(const refPlay of reversedPlays) { //this.transientDiscovered.add(refPlay); - await this.addPlayToDiscovered(refPlay); + await this.addPlayToDB(refPlay); } } } diff --git a/src/core/Api.ts b/src/core/Api.ts index e36666dc..c15976ef 100644 --- a/src/core/Api.ts +++ b/src/core/Api.ts @@ -1,6 +1,6 @@ import type { PickKeys } from "ts-essentials" import type { CompareOpKey, ComponentMinimalSelect } from "../backend/common/database/drizzle/drizzleTypes.ts" -import type { ClientType } from "./Atomic.ts" +import type { ClientType, MonitoringStatus } from "./Atomic.ts" import type { SourceType } from "./Atomic.ts" import type { ComponentType, DateLike, ErrorLike, JsonPlayObject, PlayState, QueueName, Replace, SOURCE_SOT_TYPES, SourcePlayerJson } from "./Atomic.ts" import type { Dayjs } from "dayjs" @@ -82,6 +82,7 @@ export type ComponentCommonApi = { players: Record error?: ErrorIsh warning?: ErrorIsh + monitoringStatus?: MonitoringStatus } & Omit export type ComponentCommonApiJson = Replace, string>; @@ -108,9 +109,6 @@ export type ComponentClientApiJson = Replace { return clientTypes.includes(data as ClientType); }; +export type MonitoringOrigin = 'user' | 'system'; +export const MONITORING_ORIGIN_USER: MonitoringOrigin = 'user'; +export const MONITORING_ORIGIN_SYSTEM: MonitoringOrigin = 'system'; +export interface MonitoringStatus { + monitoring: boolean + origin: MonitoringOrigin +}