diff --git a/src/backend/common/database/appMigrations/002_inputHash.ts b/src/backend/common/database/appMigrations/002_inputHash.ts new file mode 100644 index 00000000..a4afe715 --- /dev/null +++ b/src/backend/common/database/appMigrations/002_inputHash.ts @@ -0,0 +1,73 @@ +import type { SqliteDatabase, Migration } from 'sqlite-up'; +import type { MigrateBaseContext } from '../appMigrator.ts'; +import { playInputs, plays as drizzlePlays } from '../drizzle/schema/schema.ts'; +import { eq } from 'drizzle-orm'; +import { playContentBasicInvariantTransform } from '../../../utils/PlayComparisonUtils.ts'; +import { hashObject } from '../../../utils/StringUtils.ts'; + + +export const up: Migration['up'] = async (db: SqliteDatabase, ctx: MigrateBaseContext): Promise => { + + ctx.logger.info('Generating hashes for Play input data...'); + + let more = true; + let offset = 0, + processed = 0, + updated = 0; + + // first update input hashes + while (more) { + const inputRows = await ctx.db.select().from(playInputs).limit(100).offset(offset); + for (const row of inputRows) { + if (row.playHash !== null) { + processed++; + } + try { + await ctx.db.update(playInputs).set({ + playHash: hashObject(playContentBasicInvariantTransform(row.play).data) + }).where(eq(playInputs.id, row.id)); + updated++; + processed++; + } catch (e) { + ctx.logger.warn(new Error(`Failed to generate hash for Play Input ${row.id}`, { cause: e })); + } + } + offset += 100; + ctx.logger.verbose(`Play Input Hash Generation Progress: Processed ${processed} | Updated ${updated}`); + if (inputRows.length < 100) { + more = false; + } + } + + ctx.logger.info('Regenerating hashes for Play data...'); + + offset = 0; + processed = 0; + updated = 0; + more = true; + // then update Plays so hashes reflect Plays after transforms are done + while (more) { + const playsRows = await ctx.db.select().from(drizzlePlays).limit(100).offset(offset); + for (const row of playsRows) { + try { + await ctx.db.update(playInputs).set({ + playHash: hashObject(playContentBasicInvariantTransform(row.play).data) + }).where(eq(playInputs.id, row.id)); + updated++; + processed++; + } catch (e) { + ctx.logger.warn(new Error(`Failed to generate hash for Play ${row.id} (${row.uid})`, { cause: e })); + } + } + offset += 100; + ctx.logger.verbose(`Play Hash Regeneration Progress: Processed ${processed} | Updated ${updated}`); + if (playsRows.length < 100) { + more = false; + } + } +}; + +export const down: Migration['down'] = async (db: SqliteDatabase, ctx: MigrateBaseContext): Promise => { + // Rollback code here + // context is passed as ctx +}; \ No newline at end of file diff --git a/src/backend/common/database/drizzle/repositories/PlayRepository.ts b/src/backend/common/database/drizzle/repositories/PlayRepository.ts index aa513323..618db3f2 100644 --- a/src/backend/common/database/drizzle/repositories/PlayRepository.ts +++ b/src/backend/common/database/drizzle/repositories/PlayRepository.ts @@ -613,12 +613,20 @@ export class DrizzlePlayRepository extends DrizzleBaseRepository<'plays'> { return {data: res.map(x => ({...x, play: hydratePlaySelect(x, hydrate)})), meta: {limit, offset}}; } - public checkExisting = async (play: PlayObject, opts: {queueName?: string, states?: PlaySelect['state'][], taAccuracy?: TemporalAccuracy[]} & ComponentConstrainedRepoOpts = {}): Promise => { + public checkExisting = async (play: PlayObject, opts: { + queueName?: string, + states?: PlaySelect['state'][], + taAccuracy?: TemporalAccuracy[], + inputHash?: string | PlayObject, + notId?: number + } & ComponentConstrainedRepoOpts = {}): Promise => { const { queueName, componentId = this.componentId, taAccuracy = TA_DEFAULT_ACCURACY, - states + states, + inputHash, + notId, } = opts; const hash = hashObject(playContentBasicInvariantTransform(play).data); @@ -626,19 +634,25 @@ export class DrizzlePlayRepository extends DrizzleBaseRepository<'plays'> { // which we can then use with temporal comparison to make sure we are comparing the correct dates // // this isn't as fast as just comparing playDate directly but its still much faster/cheaper than paginating plays and doing everything in-memory - const dateGranularity = getTemporalAccuracyCloseVal(play.meta.source as SourceType); - let endRange: Dayjs; - if(play.data.playDateCompleted !== undefined) { - // this will be present if source reports it - // or we tracked it live with MemorySource - endRange = play.data.playDateCompleted.add(dateGranularity, 's'); - } else { - endRange = play.data.playDate.add(dateGranularity, 's'); - } + // const dateGranularity = getTemporalAccuracyCloseVal(play.meta.source as SourceType); + // let endRange: Dayjs; + // if(play.data.playDateCompleted !== undefined) { + // // this will be present if source reports it + // // or we tracked it live with MemorySource + // endRange = play.data.playDateCompleted.add(dateGranularity, 's'); + // } else { + // endRange = play.data.playDate.add(dateGranularity, 's'); + // } const where: FindWhere<'plays'> = { componentId, playedAt: buildDateCompare(getTemporallyCloseDateCompareOp(play)), }; + + if(notId !== undefined) { + where.NOT = { + id: notId + } + } if(queueName !== undefined) { where.queueStates = { @@ -653,19 +667,20 @@ export class DrizzlePlayRepository extends DrizzleBaseRepository<'plays'> { } const mbidId = playMbidIdentifier(play); - if(mbidId !== undefined) { - where.AND = [ - { - OR: [ - { - playHash: hash - }, - { - mbidIdentifier: mbidId - } - ] - } - ] + if (mbidId !== undefined || inputHash !== undefined) { + where.AND = [{ + OR: [ + { + playHash: hash + } + ] + }]; + if (mbidId !== undefined) { + where.AND[0].OR.push({ mbidIdentifier: mbidId }); + } + if (inputHash !== undefined) { + where.AND[0].OR.push({ input: { playHash: typeof inputHash === 'string' ? inputHash : hashObject(playContentBasicInvariantTransform(inputHash).data) } }); + } } else { where.playHash = hash; } @@ -673,7 +688,8 @@ export class DrizzlePlayRepository extends DrizzleBaseRepository<'plays'> { const res = await this.db.query.plays.findMany({ where, with: { - queueStates: true + queueStates: true, + input: true } }); if(res.length === 0) { diff --git a/src/backend/common/vendor/teal/TealApiClient.ts b/src/backend/common/vendor/teal/TealApiClient.ts index ea9101d4..47026807 100644 --- a/src/backend/common/vendor/teal/TealApiClient.ts +++ b/src/backend/common/vendor/teal/TealApiClient.ts @@ -1,5 +1,5 @@ import dayjs, { type Dayjs, type ManipulateType } from "dayjs"; -import type {PlayObject, PlayObjectMinimal, BrainzMeta, MBID, ScrobbleActionResult} from "../../../../core/Atomic.ts"; +import {type PlayObject, type PlayObjectMinimal, type BrainzMeta, type MBID, type ScrobbleActionResult, PARSED_FROM} from "../../../../core/Atomic.ts"; import { getRoot } from "../../../ioc.ts"; import { removeUndefinedKeys } from '../../../../core/DataUtils.ts'; import { baseFormatPlayObj } from "../../../utils/PlayTransformUtils.ts"; @@ -159,7 +159,7 @@ export const recordToPlay = (record: TealPlayRecord, options: RecordOptions = {} }, meta: { source: 'tealfm', - parsedFrom: 'history', + parsedFrom: PARSED_FROM.history, musicService, playId: options.playId, url: { diff --git a/src/backend/ioc.ts b/src/backend/ioc.ts index 4e25d402..79866d9c 100644 --- a/src/backend/ioc.ts +++ b/src/backend/ioc.ts @@ -40,6 +40,11 @@ const queuedGauge = new prom.Gauge({ help: 'Number of queued plays for a Client', labelNames: ['name', 'type'] }); +const queuedSourceGauge = new prom.Gauge({ + name: 'multiscrobbler_source_queued', + help: 'Number of queued plays for a Source', + labelNames: ['name', 'type'] +}); const deadLetterGauge = new prom.Gauge({ name: 'multiscrobbler_client_deadletter', help: 'Number of deadletter plays for a Client', @@ -144,6 +149,7 @@ const createRoot = (options: RootOptions = {logger: loggerDebug}) => { loggingConfig, sourceMetics: { discovered: discovered, + queued: queuedSourceGauge, //issues: sourceIssues }, clientMetrics: { diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index 12df1f52..0a330688 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -1416,7 +1416,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i }).catch((e) => { if (isAbortError(e)) { - const err = generateLoggableAbortReason('Dead scrrobble processing stopped', this.deadQueueAbortController.signal); + const err = generateLoggableAbortReason('Dead scrobble processing stopped', this.deadQueueAbortController.signal); this.logger.info(err); this.logger.trace(e) } else { diff --git a/src/backend/sources/AbstractSource.ts b/src/backend/sources/AbstractSource.ts index 1ed8245f..4d626bc0 100644 --- a/src/backend/sources/AbstractSource.ts +++ b/src/backend/sources/AbstractSource.ts @@ -2,7 +2,7 @@ import { 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 { type PlayMatchResult, type PlayObject, SOURCE_SOT } from "../../core/Atomic.ts"; +import { INGRESS_QUEUE, PARSED_FROM, type PlayMatchResult, type PlayObject, SOURCE_SOT } from "../../core/Atomic.ts"; import { buildTrackString, capitalize, truncateStringToLength } from "../../core/StringUtils.ts"; import AbstractComponent from "../common/AbstractComponent.ts"; import { @@ -14,7 +14,7 @@ import { type InternalConfig, type ProgressAwarePlayObject, } from "../common/infrastructure/Atomic.ts"; -import type {PlayState, PlayUserId} from '../../core/Atomic.ts'; +import type {PARSED_FROM_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"; @@ -22,6 +22,7 @@ import { TRANSFORM_HOOK } from "../../core/Transform.ts"; import TupleMap from "../common/TupleMap.ts"; import { difference, + isDebugMode, pollingBackoff, sleep, sortByOldestPlayDate, @@ -35,18 +36,20 @@ import { componentFileLogger } from '../common/logging.ts'; ; import { messageWithCausesTruncatedDefault } from "../../core/ErrorUtils.ts"; import { existingScrobble, type ExistingScrobbleOpts } from '../utils/PlayComparisonUtils.ts'; -import { staggerMapper } from '../utils/AsyncUtils.ts'; +import { consumeQueue, staggerMapper } from '../utils/AsyncUtils.ts'; import pMap, {pMapIterable} from 'p-map'; -import type { Counter } from 'prom-client'; +import type { Counter, Gauge } from 'prom-client'; import { normalizeStr } from '../utils/StringUtils.ts'; import { spawn, isAbortError, delay, throwIfAborted } from 'abort-controller-x'; -import { generateLoggableAbortReason, StageChangeError } from '../common/errors/MSErrors.ts'; +import { generateLoggableAbortReason, SimpleError, StageChangeError } from '../common/errors/MSErrors.ts'; import { DrizzlePlayRepository, playToRepositoryCreatePlayOpts, 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 { PlayWith } from '../common/database/drizzle/drizzleTypes.ts'; +import type { PlaySelect, PlaySelectWithQueueStates, PlayWith, QueueStateNew } from '../common/database/drizzle/drizzleTypes.ts'; +import { DrizzleQueueRepository } from '../common/database/drizzle/repositories/QueueRepository.ts'; +import { nanoid } from 'nanoid'; export interface RecentlyPlayedOptions { limit?: number @@ -73,12 +76,18 @@ export default abstract class AbstractSource extends AbstractComponent implement canPoll: boolean = false; polling: boolean = false; canBacklog: boolean = false; + protected discoverQueueAbortController: AbortController | undefined; + protected discoverQueuePromise: Promise | undefined; protected abortController: AbortController | undefined; protected pollingPromise: Promise | undefined; stopPollingWaitInterval: number = 200; pollRetries: number = 0; tracksDiscovered: number = 0; tracksDiscoveredTotal: number = 0; + queuedLength: number = 0; + + queueIdleMs: number = 1000; + queueConcurrency: number = 3; protected isSleeping: boolean = false; protected wakeAt: Dayjs = dayjs(); @@ -105,6 +114,9 @@ export default abstract class AbstractSource extends AbstractComponent implement declare protected componentType: 'source'; protected playRepo!: DrizzlePlayRepository; + protected queueRepo!: DrizzleQueueRepository; + + protected queuedGauge: Gauge; existingDiscoveredPlay: (playObjPre: PlayObject, existingScrobbles: PlayObject[], log?: boolean) => Promise @@ -125,7 +137,9 @@ export default abstract class AbstractSource extends AbstractComponent implement this.configDir = internal.configDir; this.emitter = emitter; - this.discoveredCounter = getRoot().items.sourceMetics.discovered; + const metrics = getRoot().items.sourceMetics; + this.discoveredCounter = metrics.discovered; + this.queuedGauge = metrics.queued; const existingScrobbleOpts: ExistingScrobbleOpts = { logger: this.logger, @@ -198,6 +212,10 @@ export default abstract class AbstractSource extends AbstractComponent implement } } if(this.isReady()) { + if(this.discoverQueuePromise === 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.'); @@ -250,6 +268,7 @@ export default abstract class AbstractSource extends AbstractComponent implement 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(); @@ -270,7 +289,9 @@ export default abstract class AbstractSource extends AbstractComponent implement protected async postDatabase(): Promise { this.playRepo = new DrizzlePlayRepository(this.db, {logger: this.logger}); + this.queueRepo = new DrizzleQueueRepository(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) { @@ -278,6 +299,13 @@ export default abstract class AbstractSource extends AbstractComponent implement } } + protected async updateQueueStats(queueNames: string[]) { + if(queueNames.includes(INGRESS_QUEUE)) { + this.queuedLength = await this.queueRepo.getQueueCount(this.dbComponent.id, [INGRESS_QUEUE]); + this.queuedGauge.labels(this.getPrometheusLabels()).set(this.queuedLength); + } + } + protected generateStaggerMappers() { const { preCompare = [], @@ -378,6 +406,62 @@ export default abstract class AbstractSource extends AbstractComponent implement // TODO make this more descriptive? or move it elsewhere recentlyPlayedTrackIsValid = (playObj: PlayObject) => true + queuePlay = async (data: PlayObject | PlayObject[]) => { + const monitoring = this.getMonitoringStatus(); + const playDatas = (Array.isArray(data) ? data : [data]).map(x => ({...x, meta: {...x.meta, wasMonitored: monitoring.monitoring, seenAt: dayjs()}})); + + const createdQueuedPlays: PlaySelect[] = []; + + await pMap(playDatas, async (queueablePlay) => { + try { + // we should be adding Plays to the queue without any transforms + // so run on "raw" play input + const cheapInputExisting = await this.playRepo.checkExisting(queueablePlay, { inputHash: queueablePlay }); + if (cheapInputExisting !== undefined) { + if(isDebugMode()) { + this.logger.trace(`Not adding ${buildTrackString(queueablePlay)} to queue because it already exists in db as Play ${cheapInputExisting.uid}`); + } + // if we have seen a play with close temporality with the exact input hash then skip it entirely + // + // this is usually the case for 'history' based plays where we are brute-force adding all plays from an api call + // and we need to prune all these duplicates + // + // but lets log if this happens and *not* history or backlog... + if(!([PARSED_FROM.history, PARSED_FROM.backlog] as PARSED_FROM_TYPE[]).includes(queueablePlay.meta.parsedFrom)) { + this.logger.warn(`Play (${buildTrackString(queueablePlay)}) dropped pre-queue due to existing (${cheapInputExisting.uid}) was not from history/backlog...`); + } + return; + } + } catch (e) { + this.logger.warn(new SimpleError('Failed to check queued scrobble for existing before adding, will continue with adding anyway', { cause: e })); + } + + // not in queue or existing queued check failed for some reason and we don't want to lose Play + const { + data, + meta + } = queueablePlay + const createPlayData = playToRepositoryCreatePlayOpts({ + play: { + data, + meta + }, + componentId: this.dbComponent.id, + state: 'queued', + }); + + const playRow = await this.playRepo.createPlays([createPlayData]); + const queueState = await this.queueRepo.create({componentId: this.dbComponent.id, playId: playRow[0].id, queueName: INGRESS_QUEUE}); + createdQueuedPlays.push(playRow[0]); + this.logger.debug(`Added ${buildTrackString(queueablePlay)} to the queue`); + this.emitPlayInsert({...playRow[0], queueStates: [queueState]} as unknown as PlayApiCommonDetailed); + this.queuedLength += 1; + this.queuedGauge.labels(this.getPrometheusLabels()).inc(); + }); + + return createdQueuedPlays; + } + protected addPlayToDB = async (play: PlayObject): Promise> => { const monitorStatus = this.getMonitoringStatus(); let state: PlayState = 'discovered'; @@ -529,25 +613,27 @@ export default abstract class AbstractSource extends AbstractComponent implement } 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}); + 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}); } - const discovered = await this.discover(backlogPlays, {discoverLocation: 'backlog', signal}); - - if (scrobbleBacklog) { - if (discovered.length > 0) { - this.logger.info('Scrobbling backlogged tracks...'); - signal.throwIfAborted(); - await this.scrobble(discovered); - this.logger.info('Backlog scrobbling complete.'); - } else { - this.logger.info('All tracks already discovered!'); - } - } else { - this.logger.info('Backlog scrobbling is disabled by config, skipping...'); - } + await this.queuePlay(backlogPlays); + this.logger.info('Backlog Plays added to discovery queue.'); + //const discovered = await this.discover(backlogPlays, {discoverLocation: 'backlog', signal}); + + // if (scrobbleBacklog) { + // if (discovered.length > 0) { + // this.logger.info('Scrobbling backlogged tracks...'); + // signal.throwIfAborted(); + // await this.scrobble(discovered); + // this.logger.info('Backlog scrobbling complete.'); + // } else { + // this.logger.info('All tracks already discovered!'); + // } + // } else { + // this.logger.info('Backlog scrobbling is disabled by config, skipping...'); + // } } return; } @@ -767,19 +853,21 @@ export default abstract class AbstractSource extends AbstractComponent implement 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); } - newDiscovered = await this.discover(playObjs, {signal}); + await this.queuePlay(playObjs); + //newDiscovered = await this.discover(playObjs, {signal}); signal.throwIfAborted(); - this.scrobble(newDiscovered, - { - forceRefresh: closeToInterval - }); + // this.scrobble(newDiscovered, + // { + // forceRefresh: closeToInterval + // }); } const activityMsgs: string[] = []; - if(newDiscovered.length > 0) { + 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) - this.lastActivityAt = newDiscovered[0].data.playDate.isAfter(this.lastActivityAt) ? newDiscovered[0].data.playDate : this.lastActivityAt; + this.lastActivityAt = playObjs[0].data.playDate.isAfter(this.lastActivityAt) ? newDiscovered[0].data.playDate : this.lastActivityAt; checkCount = 0; checksOverThreshold = 0; } @@ -841,6 +929,174 @@ export default abstract class AbstractSource extends AbstractComponent implement } } + startDiscoveryQueue = async () => { + this.setStatus('Starting discovery queue processing'); + this.discoverQueueAbortController = new AbortController(); + this.discoverQueuePromise = spawn(this.discoverQueueAbortController.signal, async (signal, { defer }) => { + await this.processDiscoveryQueue(signal); + }).catch((e) => { + const componentUpdate: Partial = { + }; + if (isAbortError(e)) { + const err = generateLoggableAbortReason('Discovery queue processing stopped', this.discoverQueueAbortController.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>(componentUpdate); + }).finally(() => { + this.discoverQueueAbortController = undefined; + this.discoverQueuePromise = undefined; + }); + } + + tryStopDiscoveryQueue = async (reason?: string | Error) => { + if(this.discoverQueuePromise === undefined) { + this.logger.verbose(`Discovery is already stopped`); + return; + } + if(this.discoverQueueAbortController === undefined) { + this.logger.error('No abort controller found! Nothing to stop.'); + return false; + } + this.discoverQueueAbortController.abort(reason) + let timePasssed = 0; + while(this.discoverQueuePromise !== 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.discoverQueuePromise !== 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); + + try { + await consumeQueue( + () => this.playRepo.getQueueNext(INGRESS_QUEUE), + 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.processQueueCurrentPlay(item, signal) + }, + { + concurrency: this.queueConcurrency, + idleMs: this.queueIdleMs, + signal, + onSuccess: () => { + taskFailures = Math.max(taskFailures - 1, 0); + }, + onError: async (e: Error) => { + 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; + } + + } + + protected processQueueCurrentPlay = async (currQueuedPlay: PlaySelectWithQueueStates, signal?: AbortSignal) => { + signal?.throwIfAborted(); + this.setStatus(`Processing Play ${currQueuedPlay.uid}`); + + const queueState = currQueuedPlay.queueStates.find(x => x.queueName === INGRESS_QUEUE); + const updatedQueueState: Partial = {}; + let state: PlayState; + try { + const preCompared = await this.transformPlay(currQueuedPlay.play, TRANSFORM_HOOK.preCompare); + let existing: PlayObject; + // cheap check for existing + const cheapExisting = await this.playRepo.checkExisting(preCompared, { notId: currQueuedPlay.id }); + if(cheapExisting !== undefined) { + updatedQueueState.error = {message: `Matched hash on existing Play ${cheapExisting.uid} with close temporality`}; + existing = {...cheapExisting.play, id: cheapExisting.id, uid: cheapExisting.uid}; + } else { + existing = await this.existingDiscovered(preCompared); + if(existing !== undefined) { + updatedQueueState.error = {message: `Matched with Play ${cheapExisting.uid}`}; + } + } + currQueuedPlay.play = preCompared; + signal?.throwIfAborted(); + if(existing === undefined) { + if(!preCompared.meta.wasMonitored) { + this.logger.debug(`Not adding ${buildTrackString(preCompared)} as discovered because monitoring was disabled when Play was created.`); + state = 'discarded'; + updatedQueueState.error = {message: 'Play was not added as discovered because monitoring was disabled when Play was created.'} + } else { + state = 'discovered'; + this.tracksDiscovered++; + this.tracksDiscoveredTotal++ + this.discoveredCounter.labels(this.getPrometheusLabels()).inc(); + this.emitEvent('discovered', {play: preCompared}); + } + } else { + this.playRepo.updateById(existing.id, {updatedAt: dayjs()}); + state = 'duped'; + currQueuedPlay.parentId = existing.id; + } + this.playRepo.updateById(currQueuedPlay.id, {play: preCompared, 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 + if(recentPlays !== undefined) { + recentPlays.push({...preCompared, id: currQueuedPlay.id, uid: currQueuedPlay.uid}); + recentPlays.sort(sortByOldestPlayDate); + this.cache.cacheDb.set(this.recentDiscoveredCacheKey(), recentPlays, '2m'); + } + updatedQueueState.queueStatus = 'completed'; + + this.logger.info(`${capitalize(state)} => ${buildTrackString(preCompared)}`); + } catch (e) { + const err = new Error(`Error ocurred while trying to discover Play ${currQueuedPlay.uid}`, {cause: e}); + updatedQueueState.error = err; + updatedQueueState.queueStatus = 'failed'; + } finally { + this.queueRepo.updateById(queueState.id, updatedQueueState); + } + + if(state === 'discovered') { + await this.scrobble([{...currQueuedPlay.play, id: currQueuedPlay.id, uid: currQueuedPlay.uid}]); + } + } + protected setIsSleeping(sleeping: boolean) { this.isSleeping = sleeping; } diff --git a/src/backend/sources/EndpointLastfmSource.ts b/src/backend/sources/EndpointLastfmSource.ts index cba73725..6d775123 100644 --- a/src/backend/sources/EndpointLastfmSource.ts +++ b/src/backend/sources/EndpointLastfmSource.ts @@ -78,10 +78,11 @@ export class EndpointLastfmSource extends MemorySource { } const discoverable = stateData.filter(x => x.play.meta.nowPlaying === false); - const discovered = await this.discover(discoverable.map(x => x.play)); - if (discovered.length > 0) { - await this.scrobble(discovered); - } + await this.queuePlay(discoverable.map(x => x.play)); + // const discovered = await this.discover(discoverable.map(x => x.play)); + // if (discovered.length > 0) { + // await this.scrobble(discovered); + // } this.componentRepo.updateById(this.dbComponent.id, {lastActiveAt: dayjs()}); this.setStatus('Waiting for Plays'); } diff --git a/src/backend/sources/EndpointListenbrainzSource.ts b/src/backend/sources/EndpointListenbrainzSource.ts index 27eebeab..37a55688 100644 --- a/src/backend/sources/EndpointListenbrainzSource.ts +++ b/src/backend/sources/EndpointListenbrainzSource.ts @@ -104,10 +104,11 @@ export class EndpointListenbrainzSource extends MemorySource { } const discoverable = stateData.filter(x => x.play.meta.nowPlaying === false && this.isValidScrobble(x.play)); - const discovered = await this.discover(discoverable.map(x => x.play)); - if (discovered.length > 0) { - await this.scrobble(discovered); - } + await this.queuePlay(discoverable.map(x => x.play)) + // const discovered = await this.discover(discoverable.map(x => x.play)); + // if (discovered.length > 0) { + // await this.scrobble(discovered); + // } this.componentRepo.updateById(this.dbComponent.id, {lastActiveAt: dayjs()}); this.setStatus('Waiting for Plays'); } diff --git a/src/backend/sources/WebScrobblerSource.ts b/src/backend/sources/WebScrobblerSource.ts index ba1ba464..43cce5e6 100644 --- a/src/backend/sources/WebScrobblerSource.ts +++ b/src/backend/sources/WebScrobblerSource.ts @@ -143,7 +143,7 @@ export class WebScrobblerSource extends MemorySource { }, meta: { trackId: uniqueID, - parsedFrom: connectorLabel, + parsedFrom: P, url: { web: trackUrl, origin: originUrl @@ -166,8 +166,8 @@ export class WebScrobblerSource extends MemorySource { return false; } - if (playObj.meta.parsedFrom !== undefined) { - const lowerSource = playObj.meta.parsedFrom.toLowerCase(); + if (playObj.meta.musicService !== undefined) { + const lowerSource = playObj.meta.musicService.toLowerCase(); if (Array.isArray(this.config.data.blacklist) && this.config.data.blacklist.length > 0) { if (this.config.data.blacklist.some(x => x === lowerSource)) { this.logger.debug(`Will not scrobble play because it is from a blacklisted connector '${lowerSource}'`); @@ -196,10 +196,11 @@ export class WebScrobblerSource extends MemorySource { this.setStatus('Received Play'); } if (stateData.play.meta.nowPlaying === false) { - const discovered = await this.discover([stateData.play]); - if (discovered.length > 0) { - await this.scrobble(discovered); - } + await this.queuePlay([stateData.play]); + // const discovered = await this.discover([stateData.play]); + // if (discovered.length > 0) { + // await this.scrobble(discovered); + // } } } diff --git a/src/backend/tests/source/TestSource.ts b/src/backend/tests/source/TestSource.ts index a7ed1082..8be9c4d7 100644 --- a/src/backend/tests/source/TestSource.ts +++ b/src/backend/tests/source/TestSource.ts @@ -4,6 +4,8 @@ import { MemoryPositionalSource } from "../../sources/MemoryPositionalSource.ts" import MemorySource from "../../sources/MemorySource.ts"; export class TestSource extends AbstractSource { + override queueIdleMs: number = 2; + override queueConcurrency: number = 1; handle(plays: PlayObject[]) { this.scrobble(plays); } diff --git a/src/backend/tests/source/source.test.ts b/src/backend/tests/source/source.test.ts index d855a608..24290151 100644 --- a/src/backend/tests/source/source.test.ts +++ b/src/backend/tests/source/source.test.ts @@ -31,6 +31,7 @@ const emitter = new WildcardEmitter(); const generateSource = async () => { const source = new TestSource('spotify', 'test-basic', {id: `test-${Date.now()}`}, {localUrl: new URL('https://example.com'), configDir: 'fake', logger: loggerTest, version: 'test'}, emitter); await source.initialize(); + await source.initTasks(); return source; } const generateMemorySource = async (config: MarkOptional = {}) => { @@ -70,7 +71,9 @@ describe('Sources use transform plays correctly', function () { const newScrobble = generatePlay({ track: 'my cool track' }); - const discovered = await source.discover([newScrobble]) + await source.queuePlay([newScrobble]); + await sleep(3); + const discovered = await source.getRecentlyDiscoveredPlays(); expect(discovered.length).eq(1); expect(discovered[0].data.track).is.eq('my fun track'); }); @@ -94,73 +97,87 @@ describe('Sources use transform plays correctly', function () { const newScrobble = generatePlay({ track: 'my cool track' }); - const discovered = await source.discover([newScrobble]) + + const pAwaiter = pEvent(source.emitter, 'discoveredToScrobble') as Promise; + + const promiseRes = await Promise.all([ + source.queuePlay([newScrobble]), + sleep(3), + pAwaiter + ]); + + //await source.queuePlay([newScrobble]); + //await sleep(3); + const discovered = await source.getRecentlyDiscoveredPlays(); expect(discovered.length).eq(1); expect(discovered[0].data.track).is.eq('my cool track'); - const pAwaiter = pEvent(source.emitter, 'discoveredToScrobble') as Promise; - source.handle(discovered); - const e = await pAwaiter; + //const pAwaiter = pEvent(source.emitter, 'discoveredToScrobble') as Promise; + //source.handle(discovered); + const e = promiseRes[2]; + expect(e).is.not.undefined; const res: PlayObject[] = !Array.isArray(e.data.data) ? [e.data.data] : e.data.data; expect(res.length).is.eq(1); expect(res[0].data.track).is.eq('my fun track'); }); - it('Transforms play existing comparison', async function() { - await using source = await generateSource(); - source.config.options = { - playTransform: { - compare: { - existing: { - type: 'user', - title: [ - { - search: 'hugely cool and very different track', - replace: 'fun' - } - ] - } - } - } - }; - source.buildTransformRules(); - const newScrobble = generatePlay({ - track: 'my hugely cool and very different track title', - }); - const discovered = await source.discover([newScrobble]) - expect(discovered.length).eq(1); - expect(discovered[0].data.track).is.eq('my hugely cool and very different track title'); - - expect((await source.discover([newScrobble])).length).is.eq(1); - }); - - it('Transforms play candidate comparison', async function() { - await using source = await generateSource(); - source.config.options = { - playTransform: { - compare: { - candidate: { - type: 'user', - title: [ - { - search: 'hugely cool and very different track', - replace: 'fun' - } - ] - } - } - } - }; - source.buildTransformRules(); - const newScrobble = generatePlay({ - track: 'my hugely cool and very different track title', - }); - const discovered = await source.discover([newScrobble]) - expect(discovered.length).eq(1); - expect(discovered[0].data.track).is.eq('my hugely cool and very different track title'); - - expect((await source.discover([newScrobble])).length).is.eq(1); - }); + // TODO need to rework these + + // it('Transforms play existing comparison', async function() { + // await using source = await generateSource(); + // source.config.options = { + // playTransform: { + // compare: { + // existing: { + // type: 'user', + // title: [ + // { + // search: 'hugely cool and very different track', + // replace: 'fun' + // } + // ] + // } + // } + // } + // }; + // source.buildTransformRules(); + // const newScrobble = generatePlay({ + // track: 'my hugely cool and very different track title', + // }); + // const discovered = await source.discover([newScrobble]) + // expect(discovered.length).eq(1); + // expect(discovered[0].data.track).is.eq('my hugely cool and very different track title'); + + // expect((await source.discover([newScrobble])).length).is.eq(1); + // }); + + // it('Transforms play candidate comparison', async function() { + // await using source = await generateSource(); + // source.config.options = { + // playTransform: { + // compare: { + // candidate: { + // type: 'user', + // title: [ + // { + // search: 'hugely cool and very different track', + // replace: 'fun' + // } + // ] + // } + // } + // } + // }; + // source.buildTransformRules(); + // const newScrobble = generatePlay({ + // track: 'my hugely cool and very different track title', + // }); + // const discovered = await source.discover([newScrobble]) + // expect(discovered.length).eq(1); + // expect(discovered[0].data.track).is.eq('my hugely cool and very different track title'); + + // expect((await source.discover([newScrobble])).length).is.eq(1); + // }); }) @@ -445,7 +462,10 @@ class DeezerTestSource extends DeezerInternalSource { const generateDeezerSource = async (options: DeezerInternalSourceOptions = {}) => { const source = new DeezerTestSource('test', {id: `test-${Date.now()}`,data: {arl: 'test'}, options}, {localUrl: new URL('https://example.com'), configDir: 'fake', logger: loggerTest, version: 'test'}, emitter); + source.queueIdleMs = 2; + source.queueConcurrency = 1; await source.initialize(); + await source.initTasks(); return source; } const firstPlayDate = dayjs().subtract(1, 'hour'); @@ -463,11 +483,13 @@ describe('Deezer Internal Source', function() { fuzzyPlay.data.playDate = targetPlay.data.playDate.add(targetPlay.data.duration, 's'); const source = await generateDeezerSource(); - source.discover([...normalizedPlays, interimPlay]); - - const discovered = await source.discover([fuzzyPlay]); - - expect(discovered.length).to.eq(1); + const queued = await source.queuePlay([...normalizedPlays, interimPlay]); + expect(queued).length(normalizedPlays.length + 1); + await Promise.race([pEvent(source.emitter, 'queueEmptied'), pEvent(source.emitter, 'discoveryQueueError')]); + expect(await source.getRecentlyDiscoveredPlays()).length(normalizedPlays.length + 1); + await source.queuePlay([fuzzyPlay]); + await Promise.race([pEvent(source.emitter, 'queueEmptied'), pEvent(source.emitter, 'discoveryQueueError')]); + expect(await source.getRecentlyDiscoveredPlays()).length(normalizedPlays.length + 2); }); }); @@ -480,11 +502,12 @@ describe('Deezer Internal Source', function() { fuzzyPlay.data.playDate = targetPlay.data.playDate.add(targetPlay.data.duration, 's'); await using source = await generateDeezerSource({fuzzyDiscoveryIgnore: true}); - await source.discover([...normalizedPlays, interimPlay]); - - const discovered = await source.discover([fuzzyPlay]); - - expect(discovered.length).to.eq(0); + const queued = await source.queuePlay([...normalizedPlays, interimPlay]); + expect(queued).length(normalizedPlays.length + 1); + await Promise.race([pEvent(source.emitter, 'queueEmptied'), pEvent(source.emitter, 'discoveryQueueError')]); + await source.queuePlay([fuzzyPlay]); + await Promise.race([pEvent(source.emitter, 'queueEmptied'), pEvent(source.emitter, 'discoveryQueueError')]); + expect(await source.getRecentlyDiscoveredPlays()).length(normalizedPlays.length + 1); }); it('discovers fuzzy play when it is the last play ', async function() { diff --git a/src/backend/utils.ts b/src/backend/utils.ts index 49a381b7..dc070294 100644 --- a/src/backend/utils.ts +++ b/src/backend/utils.ts @@ -19,12 +19,11 @@ import { NO_DEVICE } from '../core/Atomic.ts'; import type {PlayPlatformId} from '../core/Atomic.ts'; import { genGroupIdStr } from '../core/PlayUtils.ts'; import { durationToNormalizedTime } from '../core/TimeUtils.ts'; +import { setTimeout as delay } from 'node:timers/promises' dayjs.extend(utc); -export function sleep(ms: any) { - return new Promise(resolve => setTimeout(resolve, ms)); -} +export const sleep = (ms: number, opts?: Parameters[2]) => delay(ms, undefined, opts); /** sorts playObj formatted objects by playDate in ascending (oldest first) order */ export const sortByOldestPlayDate = (a: PlayObject, b: PlayObject) => { diff --git a/src/backend/utils/AsyncUtils.ts b/src/backend/utils/AsyncUtils.ts index ac22c511..b86a0141 100644 --- a/src/backend/utils/AsyncUtils.ts +++ b/src/backend/utils/AsyncUtils.ts @@ -70,4 +70,79 @@ export function staggerMapper(options: StaggerOptions) { } return await mapper(x, index); } +} + +export const consumeQueueOnce = async (next: () => Promise, process: (item: T) => Promise, opts: { + concurrency: number; + signal: AbortSignal; + onError?: (e: Error) => Promise, onSuccess?: () => void +}): Promise => { + const { concurrency, signal, onError } = opts; + signal.throwIfAborted(); + const inFlight = new Set>(); + try { + while (true) { + signal.throwIfAborted(); + if (inFlight.size >= concurrency) { + await Promise.race(inFlight); + continue; + } + const item = await next(); + if (item === undefined) break; + const task = (async () => { + try { + await process(item); + } catch (err) { + await onError?.(err); // swallow so one bad item doesn't kill the loop + } + })(); + inFlight.add(task); + void task.then(() => inFlight.delete(task)); + } + } finally { + await Promise.allSettled(inFlight); // drain before sleeping or rethrowing + } +}; + +export const consumeQueue = async ( + next: () => Promise, + process: (item: T) => Promise, + opts: { + concurrency: number; + idleMs: number; + signal: AbortSignal; + onError?: (e: Error) => Promise, + onSuccess?: () => void, + onEmpty?: () => void + }, +): Promise => { + const { concurrency, idleMs, signal, onError, onEmpty } = opts; + while (true) { + signal.throwIfAborted(); + const inFlight = new Set>(); + try { + while (true) { + signal.throwIfAborted(); + if (inFlight.size >= concurrency) { + await Promise.race(inFlight); + continue; + } + const item = await next(); + if (item === undefined) break; + const task = (async () => { + try { + await process(item); + } catch (err) { + await onError?.(err); // swallow so one bad item doesn't kill the loop + } + })(); + inFlight.add(task); + void task.then(() => inFlight.delete(task)); + } + } finally { + await Promise.allSettled(inFlight); // drain before sleeping or rethrowing + } + onEmpty?.(); + await sleep(idleMs, { signal }); + } } \ No newline at end of file diff --git a/src/client/components/ActivityDetail.tsx b/src/client/components/ActivityDetail.tsx index e60a010b..b07d61e2 100644 --- a/src/client/components/ActivityDetail.tsx +++ b/src/client/components/ActivityDetail.tsx @@ -5,7 +5,7 @@ import React, { Fragment, useEffect, useState } from "react"; import { LuChevronRight } from "react-icons/lu"; import type { MarkOptional } from "ts-essentials"; import type { ComponentsApiJson, MsSseEvent, PaginatedResponse, PlayApiCommonDetailed, QueryPlaysOptsJson, SortPlaysByProps } from "../../core/Api"; -import { CLIENT_DEAD_QUEUE, type ComponentType, type Second } from "../../core/Atomic"; +import { DEAD_QUEUE, type ComponentType, type Second } from "../../core/Atomic"; import { tanQueries, useQueryWatcher } from "../queries"; import { activityTimelineHasIssue } from "../utils/ComponentUtils"; import { ActivityTimeline } from "./ActivityTimeline"; @@ -316,7 +316,7 @@ export const ActivityStateActions = (props: {activity: PlayApiCommonDetailed}) = suffix = ; badgeProps.paddingRight = 0; } - const hasDeadQueue = queueStates.some(x => x.queueName === CLIENT_DEAD_QUEUE && x.queueStatus === 'queued'); + const hasDeadQueue = queueStates.some(x => x.queueName === DEAD_QUEUE && x.queueStatus === 'queued'); return ( diff --git a/src/client/components/ActivityTimeline.tsx b/src/client/components/ActivityTimeline.tsx index ce2a9c89..80e01d46 100644 --- a/src/client/components/ActivityTimeline.tsx +++ b/src/client/components/ActivityTimeline.tsx @@ -8,7 +8,7 @@ import { HiMiniMagnifyingGlass } from "react-icons/hi2"; import { IoMdCodeDownload } from "react-icons/io"; import { TbDatabaseEdit } from "react-icons/tb"; import type {PlayApiCommonDetailed, QueueStateApi} from "../../core/Api"; -import { CLIENT_DEAD_QUEUE, CLIENT_INGRESS_QUEUE, QUEUE_STATUS_COMPLETED, QUEUE_STATUS_FAILED, QUEUE_STATUS_QUEUED, type ComponentType, type JsonPlayObject, type LifecycleStep, type PlayMatchResult, type ScrobbleResult } from "../../core/Atomic"; +import { DEAD_QUEUE, INGRESS_QUEUE, QUEUE_STATUS_COMPLETED, QUEUE_STATUS_FAILED, QUEUE_STATUS_QUEUED, type ComponentType, type JsonPlayObject, type LifecycleStep, type PlayMatchResult, type ScrobbleResult } from "../../core/Atomic"; import { sortByNewestDate } from "../../core/PlayUtils"; import { capitalizeWords } from "../../core/StringUtils"; import { shortTodayAwareFormat } from "../../core/TimeUtils"; @@ -309,7 +309,7 @@ const QueueTimelineItem = (props: {queueState: QueueStateApi, collapsibleOpen: b - {queueState.queueName === CLIENT_DEAD_QUEUE ? 'Dead ' : ''}Queued at {shortTodayAwareFormat(dayjs(queueState.updatedAt))} + {queueState.queueName === DEAD_QUEUE ? 'Dead ' : ''}Queued at {shortTodayAwareFormat(dayjs(queueState.updatedAt))} @@ -327,7 +327,7 @@ const QueueTimelineItem = (props: {queueState: QueueStateApi, collapsibleOpen: b - {queueState.queueName === CLIENT_DEAD_QUEUE ? 'Dead ' : ''}Queue finished processing at {shortTodayAwareFormat(dayjs(queueState.updatedAt))} + {queueState.queueName === DEAD_QUEUE ? 'Dead ' : ''}Queue finished processing at {shortTodayAwareFormat(dayjs(queueState.updatedAt))} @@ -337,12 +337,12 @@ const QueueTimelineItem = (props: {queueState: QueueStateApi, collapsibleOpen: b if(queueState.queueStatus === QUEUE_STATUS_FAILED) { let titleContent: React.JSX.Element; if(queueState.error === undefined) { - titleContent = {queueState.queueName === CLIENT_DEAD_QUEUE ? 'Dead ' : ''}Queue failed at {shortTodayAwareFormat(dayjs(queueState.updatedAt))}; + titleContent = {queueState.queueName === DEAD_QUEUE ? 'Dead ' : ''}Queue failed at {shortTodayAwareFormat(dayjs(queueState.updatedAt))}; } else { titleContent = ( {queueState.queueName === CLIENT_DEAD_QUEUE ? 'Dead ' : ''}Queue failed at {shortTodayAwareFormat(dayjs(queueState.updatedAt))}} + indicator={{queueState.queueName === DEAD_QUEUE ? 'Dead ' : ''}Queue failed at {shortTodayAwareFormat(dayjs(queueState.updatedAt))}} defaultOpen={collapsibleOpen} disableUntil="md" timeline> @@ -422,7 +422,7 @@ export const ActivityTimeline = (props: ActivityTimelineProps) => { timelineItems.push(d); } - const ingressQueue = queueStates.find(x => x.queueName === CLIENT_INGRESS_QUEUE); + const ingressQueue = queueStates.find(x => x.queueName === INGRESS_QUEUE); if(ingressQueue !== undefined) { if(ingressQueue.updatedAt === ingressQueue.createdAt) { // if queue was never updated but contains extra context then only show updated @@ -436,7 +436,7 @@ export const ActivityTimeline = (props: ActivityTimelineProps) => { timelineItems.push({id: 'queue-updated-ingress', dt: dayjs(ingressQueue.updatedAt)}); } } - const deadqueue = queueStates.find(x => x.queueName === CLIENT_DEAD_QUEUE); + const deadqueue = queueStates.find(x => x.queueName === DEAD_QUEUE); if(deadqueue !== undefined) { if(deadqueue.updatedAt === deadqueue.createdAt) { // if queue was never updated but contains extra context then only show updated diff --git a/src/client/queries/index.ts b/src/client/queries/index.ts index 47703edb..ae4e9fde 100644 --- a/src/client/queries/index.ts +++ b/src/client/queries/index.ts @@ -5,7 +5,7 @@ import ky from 'ky'; import qs from 'qs'; import { baseUrl } from "../utils"; import type {ComponentsApiJson, PaginatedResponse, PlayApiCommonDetailed, PlayStateUI, QueryPlaysOptsJson} from "../../core/Api"; -import { CLIENT_DEAD_QUEUE, CLIENT_INGRESS_QUEUE, isPlayState, type SourcePlayerJson } from "../../core/Atomic"; +import { DEAD_QUEUE, INGRESS_QUEUE, isPlayState, type SourcePlayerJson } from "../../core/Atomic"; export type QueryPlaysOptsJsonRefreshable = Omit & {nonce?: string, state?: PlayStateUI[]}; @@ -42,12 +42,12 @@ const activities = createQueryKeys('activities', { derived.state = state.filter(x => isPlayState(x)); // remove 'dead queued' derived play state and replace with filter for queue = 'dead' & state = 'queued' - if(state.includes('dead queued') && !rest.queues?.some(x => x.queueName === CLIENT_DEAD_QUEUE)) { - derived.queues = [...(rest.queues ?? []), {queueName: CLIENT_DEAD_QUEUE, queueStatus: 'queued'}]; + if(state.includes('dead queued') && !rest.queues?.some(x => x.queueName === DEAD_QUEUE)) { + derived.queues = [...(rest.queues ?? []), {queueName: DEAD_QUEUE, queueStatus: 'queued'}]; } // remove 'queued' play state and replace with filter for queue = 'ingress' & state = 'queued' - if(state.includes('queued') && !rest.queues?.some(x => x.queueName === CLIENT_INGRESS_QUEUE)) { - derived.queues = [...(derived.queues ?? []), {queueName: CLIENT_INGRESS_QUEUE, queueStatus: 'queued'}]; + if(state.includes('queued') && !rest.queues?.some(x => x.queueName === INGRESS_QUEUE)) { + derived.queues = [...(derived.queues ?? []), {queueName: INGRESS_QUEUE, queueStatus: 'queued'}]; derived.state = derived.state.filter(x => x !== 'queued'); } } diff --git a/src/core/Atomic.ts b/src/core/Atomic.ts index b3cf4c66..939f058b 100644 --- a/src/core/Atomic.ts +++ b/src/core/Atomic.ts @@ -201,9 +201,9 @@ export interface PlayMetaBase { musicService?: string /** - * Specifies from what facet/data from the source this play was parsed from IE history, now playing, etc... + * Specifies from what facet/data from the source this play was parsed from IE player, backlog, now playing, etc... * */ - parsedFrom?: string + parsedFrom?: PARSED_FROM_TYPE /** * Unique ID for this track, given by the Source * */ @@ -267,6 +267,9 @@ export interface PlayMetaBase { comment?: string + /** Was the component activitely monitoring when this Play was created? */ + wasMonitored?: boolean + //lifecycle: PlayLifecycle lifecycleInputs?: LifecycleInput[] @@ -477,6 +480,14 @@ export const SOURCE_SOT = { } as const satisfies Record export const sourceSotTypes: SOURCE_SOT_TYPES[] = ['player','history','ingress']; +export type PARSED_FROM_TYPE = 'backlog' | 'now playing' | 'player' | 'history'; +export const PARSED_FROM = { + backlog : 'backlog', + nowPlaying: 'now playing', + player: 'player', + history: 'history' +} as const satisfies Record + export interface URLData { url: URL normal: string