diff --git a/src/backend/common/database/drizzle/entityUtils.ts b/src/backend/common/database/drizzle/entityUtils.ts index 84a241a6..67a2e6ad 100644 --- a/src/backend/common/database/drizzle/entityUtils.ts +++ b/src/backend/common/database/drizzle/entityUtils.ts @@ -149,4 +149,7 @@ export const scrobbleToPlayEvent = (data: PlayEventScrobbleResultData): Omit `play` in obj && typeof obj.play === 'object' + && `componentId` in obj && typeof obj.componentId === 'number' \ No newline at end of file diff --git a/src/backend/sources/AbstractSource.ts b/src/backend/sources/AbstractSource.ts index b7bcf02d..8013ba7f 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 { INGRESS_QUEUE, PARSED_FROM, type PlayMatchResult, type PlayObject, QUEUE_STATUS_COMPLETED, SOURCE_SOT } from "../../core/Atomic.ts"; +import { INGRESS_QUEUE, isPlayObject, PARSED_FROM, PLAY_STATES, type PlayMatchResult, type PlayObject, QUEUE_STATUS_COMPLETED, 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 {PARSED_FROM_TYPE, PlayState, PlayUserId} from '../../core/Atomic.ts'; +import type {PARSED_FROM_TYPE, PlayState, 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"; @@ -47,11 +47,11 @@ 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 { PlaySelect, PlaySelectWithQueueStates, QueueStateNew, QueueStateSelect } from '../common/database/drizzle/drizzleTypes.ts'; +import type { PlayEventSelect, PlaySelect, PlaySelectWithQueueStates, PlayWith, QueueStateNew, QueueStateSelect } from '../common/database/drizzle/drizzleTypes.ts'; import { DrizzleQueueRepository } from '../common/database/drizzle/repositories/QueueRepository.ts'; import { DrizzlePlayEventsRepository } from '../common/database/drizzle/repositories/PlayEventsRepository.ts'; import { PLAY_EVENT_TYPE, type PlayEvent } from '../../core/PlayEvent.ts'; -import { dupeCheckToPlayEvent, queueStateToPlayEvent, stateChangeToPlayEvent, transformToPlayEvent } from '../common/database/drizzle/entityUtils.ts'; +import { dupeCheckToPlayEvent, entityIsPlayEntity, queueStateToPlayEvent, stateChangeToPlayEvent, transformToPlayEvent } from '../common/database/drizzle/entityUtils.ts'; export interface RecentlyPlayedOptions { limit?: number @@ -375,70 +375,95 @@ 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()}})); - + queuePlay = async (data: (PlayObject | PlayObject[]) | (PlaySelectWithQueueStates | PlaySelectWithQueueStates[]), context?: QueueContext) => { const createdQueuedPlays: PlaySelect[] = []; - await pMap(playDatas, async (queueablePlay) => { - try { - // 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. - 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 }); - if (cheapInputExisting !== undefined) { - if (isDebugMode() || PARSED_FROM.backlog === queueablePlay.meta.parsedFrom) { - // 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}`); + const dataArray = Array.isArray(data) ? data : [data]; + if (dataArray.every(x => entityIsPlayEntity(x))) { + for (const playSelect of dataArray) { + let queue = playSelect.queueStates.find(x => x.queueName === INGRESS_QUEUE); + if (queue === undefined) { + queue = await this.queueRepo.create({ componentId: this.dbComponent.id, playId: playSelect.id, queueName: INGRESS_QUEUE, context }) as QueueStateSelect; + } else { + this.queueRepo.updateById(queue.id, { queueStatus: 'queued', context }); + } + const events = await this.playEventsRepo.createMany([ + { playId: playSelect.id, ...stateChangeToPlayEvent({ state: 'queued' }) }, + { playId: playSelect.id, ...queueStateToPlayEvent({ ...queue, context: context ?? queue.context }) } + ]) as PlayEventSelect[]; + playSelect.state = 'queued'; + await this.playRepo.updateById(playSelect.id, {state: 'queued'}); + if (`events` in playSelect) { + (playSelect as PlayWith<'events'>).events = (playSelect as PlayWith<'events'>).events.concat(events); + } else { + (playSelect as unknown as PlayWith<'events'>).events = events; + } + this.emitPlayUpdate({ ...playSelect } as unknown as PlayApiCommonDetailed); + } + } else if (dataArray.every(x => isPlayObject(x))) { + const monitoring = this.getMonitoringStatus(); + const playDatas = dataArray.map(x => ({ ...x, meta: { ...x.meta, wasMonitored: monitoring.monitoring, seenAt: dayjs() } })); + + await pMap(playDatas, async (queueablePlay) => { + try { + // 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. + 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 }); + if (cheapInputExisting !== undefined) { + if (isDebugMode() || PARSED_FROM.backlog === queueablePlay.meta.parsedFrom) { + // 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; } - return; } + } catch (e) { + this.logger.warn(new SimpleError('Failed to check queued scrobble for existing before adding, will continue with adding anyway', { cause: e })); } - } 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, - original - } = queueablePlay - const createPlayData = playToRepositoryCreatePlayOpts({ - play: { + + // not in queue or existing queued check failed for some reason and we don't want to lose Play + const { data, meta, original - }, - componentId: this.dbComponent.id, - state: 'queued', + } = queueablePlay + const createPlayData = playToRepositoryCreatePlayOpts({ + play: { + data, + meta, + original + }, + 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 }) as QueueStateSelect; + await this.playEventsRepo.createMany([ + { playId: playRow[0].id, ...stateChangeToPlayEvent({ state: 'queued' }), createdAt: playRow[0].seenAt.add(1, 'ms') }, + { playId: playRow[0].id, ...queueStateToPlayEvent(queueState), createdAt: queueState.createdAt } + ]); + 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(); }); - - const playRow = await this.playRepo.createPlays([createPlayData]); - const queueState = await this.queueRepo.create({componentId: this.dbComponent.id, playId: playRow[0].id, queueName: INGRESS_QUEUE}) as QueueStateSelect; - await this.playEventsRepo.createMany([ - {playId: playRow[0].id, ...stateChangeToPlayEvent({state: 'queued'}), createdAt: playRow[0].seenAt.add(1,'ms')}, - {playId: playRow[0].id, ...queueStateToPlayEvent(queueState), createdAt: queueState.createdAt} - ]); - 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(); - }); - + } else { + throw new Error('Data passed to queuePlay must be either all be PlayObject or all PlaySelect objects'); + } return createdQueuedPlays; } @@ -983,7 +1008,7 @@ export default abstract class AbstractSource extends AbstractComponent implement let events: Omit[] = []; try { - if(!currQueuedPlay.play.meta.wasMonitored) { + if(queueState.context?.isRetry !== true && !currQueuedPlay.play.meta.wasMonitored) { this.logger.debug(`Not processing ${buildTrackString(currQueuedPlay.play)} because monitoring was disabled when Play was queued.`); state = 'discarded'; events.push(stateChangeToPlayEvent({state, reason: 'Not processing because monitoring was disabled when Play was queued'})); diff --git a/src/core/Atomic.ts b/src/core/Atomic.ts index 1ff05011..c1eab7f1 100644 --- a/src/core/Atomic.ts +++ b/src/core/Atomic.ts @@ -659,6 +659,7 @@ export interface QueueContext { transform?: boolean dupeCheck?: boolean useCache?: boolean + isRetry?: boolean } /** diff --git a/src/core/PlayEvent.ts b/src/core/PlayEvent.ts index c4ec05e7..c15f8e1a 100644 --- a/src/core/PlayEvent.ts +++ b/src/core/PlayEvent.ts @@ -1,5 +1,5 @@ import type { Dayjs } from "dayjs"; -import type { DateLike, ErrorLike, LifecycleStep, PlayMatchResult, PlayState, QueueStatus, ScrobbleResult } from "./Atomic.ts"; +import type { DateLike, ErrorLike, LifecycleStep, PlayMatchResult, PlayState, QueueContext, QueueStatus, ScrobbleResult } from "./Atomic.ts"; export type PlayEventType = 'transform' | 'queueStateChange' | 'playStateChange' | 'dupeCheck' | 'scrobbleResult'; export const PLAY_EVENT_TYPE = { @@ -25,6 +25,7 @@ export type PlayEventTransform = BasePlayEvent<'tran export interface PlayEventQueueStateChangeData { queueName: string queueStatus: QueueStatus + context?: QueueContext retries?: number error?: ErrorLike }