diff --git a/src/backend/common/database/drizzle/entityUtils.ts b/src/backend/common/database/drizzle/entityUtils.ts index d7cec546..84a241a6 100644 --- a/src/backend/common/database/drizzle/entityUtils.ts +++ b/src/backend/common/database/drizzle/entityUtils.ts @@ -3,13 +3,13 @@ import type {PlayHistoricalNew, PlayHistoricalSelect, PlayNew, PlaySelect, PlayS import type {PlayInputNew} from "./drizzleTypes.ts"; import type {QueueStateNew} from "./drizzleTypes.ts"; import type {ComponentNew} from "./drizzleTypes.ts"; -import type { MarkOptional } from "ts-essentials"; -import { DEAD_QUEUE, type DeadLetterScrobble, type ErrorLike, type PlayObject } from "../../../../core/Atomic.ts"; +import type { MarkOptional, MarkRequired } from "ts-essentials"; +import { DEAD_QUEUE, type DeadLetterScrobble, type ErrorLike, type LifecycleStep, type PlayObject } from "../../../../core/Atomic.ts"; import dayjs from "dayjs"; import { playContentBasicInvariantTransform, playMbidIdentifier } from "../../../utils/PlayComparisonUtils.ts"; import { hashObject } from "../../../utils/StringUtils.ts"; import { serializeError } from "serialize-error"; -import type { PlayEventQueueStateChange, PlayEventQueueStateChangeData } from "../../../../core/PlayEvent.ts"; +import { PLAY_EVENT_TYPE, type PlayEventDupeCheck, type PlayEventDupeCheckData, type PlayEventPlayStateChange, type PlayEventPlayStateChangeData, type PlayEventQueueStateChange, type PlayEventQueueStateChangeData, type PlayEventScrobbleResult, type PlayEventScrobbleResultData, type PlayEventTransform } from "../../../../core/PlayEvent.ts"; export const generateComponentEntity = (data: MarkOptional): ComponentNew => { assert(data.name !== undefined, 'Must provide name'); @@ -109,4 +109,44 @@ export const queueStateToEventData = (qs: QueueStateSelect): PlayEventQueueState error, retries } -} \ No newline at end of file +} + +export const transformToPlayEvent = (lifecycle: LifecycleStep[]): Omit => ({ + eventName: PLAY_EVENT_TYPE.transform, + createdAt: dayjs(lifecycle[0].createdAt), + data: lifecycle +}); + +export const stateChangeToPlayEvent = (partial: PlayEventPlayStateChangeData): Omit => ({ + eventName: PLAY_EVENT_TYPE.playStateChange, + createdAt: dayjs(), + data: partial +}) + +export const queueStateToPlayEvent = (partial: QueueStateSelect): Omit => ({ + eventName: PLAY_EVENT_TYPE.queueStateChange, + createdAt: dayjs(), + data: partial +}); + +export const dupeCheckToPlayEvent = (partial: MarkRequired, 'match'>): Omit => { + const {match, ...rest} = partial; + const data: PlayEventDupeCheckData = { + match, + score: match ? 1 : 0, + breakdowns: [], + createdAt: dayjs().toISOString(), + ...rest + }; + return { + eventName: PLAY_EVENT_TYPE.dupeCheck, + createdAt: dayjs(), + data + } +} + +export const scrobbleToPlayEvent = (data: PlayEventScrobbleResultData): Omit => ({ + eventName: PLAY_EVENT_TYPE.scrobbleResult, + createdAt: data.createdAt ?? dayjs(), + data +}); \ No newline at end of file diff --git a/src/backend/common/database/drizzle/repositories/BaseRepository.ts b/src/backend/common/database/drizzle/repositories/BaseRepository.ts index 8613f4b8..9616f212 100644 --- a/src/backend/common/database/drizzle/repositories/BaseRepository.ts +++ b/src/backend/common/database/drizzle/repositories/BaseRepository.ts @@ -45,9 +45,9 @@ export abstract class DrizzleBaseRepository { await this.db.delete(this.table).where(inArray(this.table.id, ids)); } - async updateById(id: number, data: Partial): Promise { + async updateById(id: number, data: Partial): Promise { assert(id !== null && id !== undefined, `${id === null ? 'null' : 'undefined'} given for entity id`); - await this.db.update(this.table).set(data).where(eq(this.table.id, id)); + return (await this.db.update(this.table).set(data).where(eq(this.table.id, id)).returning())[0]; } async create(data: typeof this.table.$inferInsert): Promise { diff --git a/src/backend/common/database/drizzle/repositories/PlayRepository.ts b/src/backend/common/database/drizzle/repositories/PlayRepository.ts index 460276b9..207ef550 100644 --- a/src/backend/common/database/drizzle/repositories/PlayRepository.ts +++ b/src/backend/common/database/drizzle/repositories/PlayRepository.ts @@ -10,15 +10,15 @@ import { playContentBasicInvariantTransform, playMbidIdentifier } from "../../.. import { hashObject } from "../../../../utils/StringUtils.ts"; import { comparePlayTemporally, getScrobbleTsSOCDateWithContext, getTemporalAccuracyCloseVal, hasAcceptableTemporalAccuracy } from "../../../../utils/TimeUtils.ts"; import { type CompactableProperty, type RetentionOptions, retentionPlayTypes } from "../../../infrastructure/config/database.ts"; -import type {SourceType} from "../../../../../core/Atomic.ts"; +import type {ErrorLike, SourceType} from "../../../../../core/Atomic.ts"; import type {FindMany, FindWhere, FindWith, PlayInputNew, PlayNew, PlaySelect, PlaySelectWithQueueStates, PlayWith, QueueStateSelect, WhereClause} from "../drizzleTypes.ts"; import { type DbConcrete, runTransaction } from "../drizzleUtils.ts"; -import { generateInputEntity, generatePlayEntity, hydratePlaySelect, type PlayEntityOpts, type PlayHydateOptions } from "../entityUtils.ts"; +import { generateInputEntity, generatePlayEntity, hydratePlaySelect, stateChangeToPlayEvent, transformToPlayEvent, type PlayEntityOpts, type PlayHydateOptions } from "../entityUtils.ts"; import { playEvents, playInputs, plays, relations } from "../schema/schema.ts"; import { buildDateCompare, type CompareDateOp, type ComponentConstrainedRepoOpts, DrizzleBaseRepository, type DrizzleRepositoryOpts } from "./BaseRepository.ts"; import type {PaginatedResponse} from "../../../../../core/Api.ts"; import type {PaginatedQueryResponse} from "../../../../../core/Api.ts"; -import type { PlayEventTransform } from "../../../../../core/PlayEvent.ts"; +import { type PlayEventTransform } from "../../../../../core/PlayEvent.ts"; // https://github.com/drizzle-team/drizzle-orm/issues/695 may be useful for typing models with relations? @@ -138,13 +138,7 @@ export class DrizzlePlayRepository extends DrizzleBaseRepository<'plays'> { } = entitiesOpts[index]; if(lifecycle.length > 0) { - const transformEvent: PlayEventTransform = { - playId: x.id, - eventName: 'transform', - createdAt: dayjs(lifecycle[0].createdAt), - data: lifecycle - } - return transformEvent; + return {...transformToPlayEvent(lifecycle), playId: x.id} } return undefined; }).filter(x => x !== undefined); @@ -777,6 +771,20 @@ where compacted IS NOT NULL group by componentId,compacted;`); return res; } + + async updateById(id: number, data: Partial & {event?: boolean, reason?: string, error?: ErrorLike}): Promise { + const res = await super.updateById(id, data) as PlaySelect; + if(data.event === true) { + if(data.state !== undefined) { + try { + await this.db.insert(playEvents).values({...stateChangeToPlayEvent(removeUndefinedKeys({state: data.state, reason: data.reason, error: data.error})), playId: id}); + } catch (e) { + this.logger.warn(new Error(`Failed to create Play Event for state change ${data.state} on Play ${id}`)); + } + } + } + return res; + } } export const getTemporallyCloseDateCompareOp = (play: PlayObject, opts: {bufferTime?: number, useCompleted?: boolean, useDuration?: boolean} = {}): CompareDateOp => { diff --git a/src/backend/common/database/drizzle/repositories/QueueRepository.ts b/src/backend/common/database/drizzle/repositories/QueueRepository.ts index 52e53e78..91b8da6f 100644 --- a/src/backend/common/database/drizzle/repositories/QueueRepository.ts +++ b/src/backend/common/database/drizzle/repositories/QueueRepository.ts @@ -1,9 +1,10 @@ import { eq, and, lte, inArray } from "drizzle-orm"; import { DrizzleBaseRepository, type DrizzleRepositoryOpts } from "./BaseRepository.ts"; import type {DbConcrete} from "../drizzleUtils.ts"; -import type {QueueStateSelect} from "../drizzleTypes.ts"; -import { queueStates } from "../schema/schema.ts"; +import type {PlaySelect, QueueStateSelect} from "../drizzleTypes.ts"; +import { playEvents, queueStates } from "../schema/schema.ts"; import { DEAD_QUEUE } from "../../../../../core/Atomic.ts"; +import { queueStateToPlayEvent } from "../entityUtils.ts"; export class DrizzleQueueRepository extends DrizzleBaseRepository<'queueStates'> { constructor(db: DbConcrete, opts: DrizzleRepositoryOpts = {}) { @@ -38,4 +39,28 @@ export class DrizzleQueueRepository extends DrizzleBaseRepository<'queueStates'> inArray(queueStates.queueStatus, queueStatus) )); } + + async create(data: typeof this.table.$inferInsert & {playId?: PlaySelect['id']}): Promise { + const res = await super.create(data) as QueueStateSelect; + if(data.playId !== undefined) { + try { + await this.db.insert(playEvents).values({...queueStateToPlayEvent(res), playId: data.playId}); + } catch (e) { + this.logger.warn(new Error(`Failed to create Play Event for new queue creation on Play ${data.playId}`)); + } + } + return res; + } + + async updateById(id: number, data: Partial & {playId?: PlaySelect['id']}): Promise { + const res = await super.updateById(id, data) as QueueStateSelect; + if(data.playId !== undefined) { + try { + await this.db.insert(playEvents).values({...queueStateToPlayEvent(res), playId: data.playId}); + } catch (e) { + this.logger.warn(new Error(`Failed to create Play Event for queue ${res.queueName} on Play ${data.playId}`)); + } + } + return res; + } } \ No newline at end of file diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index 5d62c7ab..2c6fb13b 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -74,8 +74,8 @@ import assert from "node:assert"; import { COMPONENT_STATE, type ComponentClientApiJson, type PlayApiCommonDetailed } from "../../core/Api.ts"; import type {ComponentState} from "react"; import { DrizzlePlayEventsRepository } from "../common/database/drizzle/repositories/PlayEventsRepository.ts"; -import { PLAY_EVENT_TYPE, type PlayEvent } from "../../core/PlayEvent.ts"; -import { queueStateToEventData } from "../common/database/drizzle/entityUtils.ts"; +import { type PlayEvent } from "../../core/PlayEvent.ts"; +import { dupeCheckToPlayEvent, queueStateToPlayEvent, scrobbleToPlayEvent, stateChangeToPlayEvent, transformToPlayEvent } from "../common/database/drizzle/entityUtils.ts"; type SourceMappedPlayer = {player: SourcePlayerObj, source: SourceIdentifier}; type PlatformMappedPlays = Map; @@ -1269,7 +1269,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i } if (historicalError === undefined) { const { summary, ...matchResult } = await this.existingScrobble(currQueuedPlay.play, historicalPlays); - events.push({eventName: PLAY_EVENT_TYPE.dupeCheck, data: {summary, ...matchResult}, createdAt: dayjs()}); + events.push(dupeCheckToPlayEvent({summary, ...matchResult})); // currQueuedPlay.play.scrobble = { // ...(currQueuedPlay.play.scrobble ?? {}), // match: matchResult, @@ -1281,13 +1281,13 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i const { lifecycle = [] } = transformedScrobble; const psLifecycle = lifecycle.filter(x => x.hook === TRANSFORM_HOOK.postCompare); if(psLifecycle.length > 0) { - events.push({eventName: PLAY_EVENT_TYPE.transform, createdAt: dayjs(psLifecycle[0].createdAt), data: psLifecycle}); + events.push(transformToPlayEvent(psLifecycle)); } signal.throwIfAborted(); try { const scrobbledPlay = await this.scrobble(transformedScrobble, {signal}); const {scrobble} = scrobbledPlay; - events.push({eventName: PLAY_EVENT_TYPE.scrobbleResult, createdAt: scrobble.createdAt, data: scrobble}); + events.push(scrobbleToPlayEvent(scrobble)); //currQueuedPlay.play = scrobbledPlay; await this.addScrobbledTrack(scrobbledPlay); //handledShiftedPlay = true; @@ -1305,7 +1305,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i scrobbleRes.payload = this.playToClientPayload(transformedScrobble); scrobbleRes.error = serializeError(e); } - events.push({eventName: PLAY_EVENT_TYPE.scrobbleResult, createdAt: scrobbleRes.createdAt, data: scrobbleRes}); + events.push(scrobbleToPlayEvent(scrobbleRes)); queueError = e; deadQueueEntity = await this.addDeadLetterScrobble(currQueuedPlay, e); //handledShiftedPlay = true; @@ -1350,17 +1350,17 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i queueState.queueStatus = 'failed'; queueState.error = queueError; await this.playRepo.updateById(currQueuedPlay.id, {state: 'failed', error: queueError, play: currQueuedPlay.play}); - events.push({eventName: PLAY_EVENT_TYPE.playStateChange, data: {state: 'failed'}, createdAt: dayjs()}); - events.push({eventName: PLAY_EVENT_TYPE.queueStateChange, data: queueStateToEventData(queueState), createdAt: dayjs()}); + events.push(stateChangeToPlayEvent({state: 'failed'})); + events.push(queueStateToPlayEvent(queueState)); currQueuedPlay.state = 'failed'; //currQueuedPlay.error = queueError; } else { await this.queueRepo.updateById(queueState.id, {queueStatus: 'completed'}); await this.playRepo.updateById(currQueuedPlay.id, {state: successState ?? 'scrobbled', play: currQueuedPlay.play}); - events.push({eventName: PLAY_EVENT_TYPE.playStateChange, data: {state: successState ?? 'scrobbled'}, createdAt: dayjs()}); + events.push(stateChangeToPlayEvent({state: successState ?? 'scrobbled'})); currQueuedPlay.state = successState ?? 'scrobbled'; queueState.queueStatus = 'completed'; - events.push({eventName: PLAY_EVENT_TYPE.queueStateChange, data: queueStateToEventData(queueState), createdAt: dayjs()}); + events.push(queueStateToPlayEvent(queueState)); } await this.playEventsRepo.createMany(events.map(x => ({...x, playId: currQueuedPlay.id}))); this.emitPlayUpdate({...currQueuedPlay, queueStates: [queueState]} as unknown as PlayApiCommonDetailed); @@ -1498,7 +1498,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i this.deadLogger.warn(new SimpleError(`${deadScrobble.uid} - ${buildTrackString(deadScrobble.play)} from Source '${deadScrobble.play.meta.source}' => cannot get historical scrobbles`, { cause: e, shortStack: true })); } - events.push({eventName: PLAY_EVENT_TYPE.queueStateChange, data: queueStateToEventData({...deadQueueState, queueStatus: 'failed', error: e}), createdAt: dayjs()}); + events.push(queueStateToPlayEvent({...deadQueueState, queueStatus: 'failed', error: e})); this.queueRepo.updateById(deadQueueState.id, { retries: deadQueueState.retries + 1, error: e, updatedAt: dayjs(), queueStatus: 'failed' }); //this.playRepo.updateById(deadScrobble.id, {error: e}); // deadScrobble.retries++; @@ -1511,7 +1511,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i } signal?.throwIfAborted(); const { summary, ...matchResult } = await this.existingScrobble(deadScrobble.play, historicalPlays); - events.push({eventName: PLAY_EVENT_TYPE.dupeCheck, data: {summary, ...matchResult}, createdAt: dayjs()}); + events.push(dupeCheckToPlayEvent({summary, ...matchResult})) // deadScrobble.play.scrobble = { // ...(deadScrobble.play.scrobble ?? {}), // match: matchResult, @@ -1522,19 +1522,19 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i const { lifecycle = [] } = transformedScrobble; const psLifecycle = lifecycle.filter(x => x.hook === TRANSFORM_HOOK.postCompare); if(psLifecycle.length > 0) { - events.push({eventName: PLAY_EVENT_TYPE.transform, createdAt: dayjs(psLifecycle[0].createdAt), data: psLifecycle}); + events.push(transformToPlayEvent(psLifecycle)); } signal?.throwIfAborted(); try { const scrobbledPlay = await this.scrobble(transformedScrobble); const {scrobble} = scrobbledPlay; - events.push({eventName: PLAY_EVENT_TYPE.scrobbleResult, createdAt: scrobble.createdAt, data: scrobble}); + events.push(scrobbleToPlayEvent(scrobble)); deadScrobble.play = scrobbledPlay; await this.addScrobbledTrack(scrobbledPlay); this.playRepo.updateById(deadScrobble.id, { play: deadScrobble.play, state: 'scrobbled' }); this.queueRepo.updateById(deadQueueState.id, { error: null, updatedAt: dayjs(), queueStatus: QUEUE_STATUS_COMPLETED }); - events.push({eventName: PLAY_EVENT_TYPE.queueStateChange, data: {queueName: DEAD_QUEUE, queueStatus: QUEUE_STATUS_COMPLETED}, createdAt: dayjs()}); - events.push({eventName: PLAY_EVENT_TYPE.playStateChange, data: {state: 'scrobbled'}, createdAt: dayjs()}); + events.push(queueStateToPlayEvent({...deadQueueState, queueStatus: QUEUE_STATUS_COMPLETED})); + events.push(stateChangeToPlayEvent({state: 'scrobbled'})); this.removeDeadLetterScrobble(deadScrobble, 'scrobbled', true); } catch (e) { const scrobbleRes: ScrobbleResult = { @@ -1553,9 +1553,9 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i scrobbleRes.payload = this.playToClientPayload(transformedScrobble); scrobbleRes.error = serializeError(e); } - events.push({eventName: PLAY_EVENT_TYPE.scrobbleResult, createdAt: scrobbleRes.createdAt, data: scrobbleRes}); + events.push(scrobbleToPlayEvent(scrobbleRes)); this.queueRepo.updateById(deadQueueState.id, { retries: deadQueueState.retries + 1, error: e, updatedAt: dayjs(), queueStatus: 'failed' }); - events.push({eventName: PLAY_EVENT_TYPE.queueStateChange, data: queueStateToEventData({...deadQueueState, queueStatus: 'failed', error: e}), createdAt: dayjs()}); + events.push(queueStateToPlayEvent({...deadQueueState, queueStatus: 'failed', error: e})); //this.playRepo.updateById(deadScrobble.id, { play: deadScrobble.play }); // deadScrobble.retries++; // deadScrobble.error = messageWithCauses(e); @@ -1569,14 +1569,14 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i this.playRepo.updateById(deadScrobble.id, { play: deadScrobble.play }); this.deadLogger.verbose(`Looks like ${buildTrackString(deadScrobble.play)} was already scrobbled!\n${summary}`); this.removeDeadLetterScrobble(deadScrobble, 'duped', true); - events.push({eventName: PLAY_EVENT_TYPE.queueStateChange, data: {queueName: DEAD_QUEUE, queueStatus: QUEUE_STATUS_COMPLETED}, createdAt: dayjs()}); - events.push({eventName: PLAY_EVENT_TYPE.playStateChange, data: {state: 'duped', reason: 'Looks like it was already scrobbled downstream'}, createdAt: dayjs()}); + events.push(queueStateToPlayEvent({...deadQueueState, queueStatus: QUEUE_STATUS_COMPLETED})); + events.push(stateChangeToPlayEvent({state: 'duped', reason: 'Looks like it was already scrobbled downstream'})); } return [true, deadScrobble]; } catch (e) { if(deadQueueState !== undefined) { - events.push({eventName: PLAY_EVENT_TYPE.queueStateChange, data: queueStateToEventData({...deadQueueState, queueStatus: 'failed', error: e}), createdAt: dayjs()}); + events.push(queueStateToPlayEvent({...deadQueueState, queueStatus: 'failed', error: e})); } } finally { await this.playEventsRepo.createMany(events.map(x => ({...x, playId: deadScrobble.id}))); @@ -1701,8 +1701,8 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i 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, eventName: PLAY_EVENT_TYPE.playStateChange, data: {state: 'queued'}, createdAt: playRow[0].seenAt.add(1,'ms')}, - {playId: playRow[0].id, eventName: PLAY_EVENT_TYPE.queueStateChange, data: queueStateToEventData(queueState), createdAt: queueState.createdAt} + {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(play)} to the queue`); @@ -1738,7 +1738,7 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i queueName: DEAD_QUEUE }) as QueueStateSelect; await this.playEventsRepo.createMany([ - {playId: data.id, eventName: PLAY_EVENT_TYPE.queueStateChange, data: queueStateToEventData(newQueue), createdAt: newQueue.createdAt} + {playId: data.id, ...queueStateToPlayEvent(newQueue), createdAt: newQueue.createdAt} ]); const deadData = {id: nanoid(), retries: 0, error: e, play: data.play}; //this.deadLetterScrobbles.push(deadData); diff --git a/src/backend/sources/AbstractSource.ts b/src/backend/sources/AbstractSource.ts index 75e130d9..5c8ff034 100644 --- a/src/backend/sources/AbstractSource.ts +++ b/src/backend/sources/AbstractSource.ts @@ -50,8 +50,8 @@ import type {PaginatedResponse} from "../../core/Api.ts"; import type { PlaySelect, PlaySelectWithQueueStates, 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 { queueStateToEventData } from '../common/database/drizzle/entityUtils.ts'; +import { type PlayEvent } from '../../core/PlayEvent.ts'; +import { dupeCheckToPlayEvent, queueStateToPlayEvent, stateChangeToPlayEvent, transformToPlayEvent } from '../common/database/drizzle/entityUtils.ts'; export interface RecentlyPlayedOptions { limit?: number @@ -427,8 +427,8 @@ export default abstract class AbstractSource extends AbstractComponent implement 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, eventName: PLAY_EVENT_TYPE.playStateChange, data: {state: 'queued'}, createdAt: playRow[0].seenAt.add(1,'ms')}, - {playId: playRow[0].id, eventName: PLAY_EVENT_TYPE.queueStateChange, data: queueStateToEventData(queueState), createdAt: queueState.createdAt} + {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`); @@ -515,7 +515,7 @@ export default abstract class AbstractSource extends AbstractComponent implement const {lifecycle = []} = p; const psLifecycle = lifecycle.filter(x => x.hook === TRANSFORM_HOOK.postCompare); if(psLifecycle.length > 0) { - events.push({playId: p.id, eventName: PLAY_EVENT_TYPE.transform, createdAt: dayjs(psLifecycle[0].createdAt), data: psLifecycle}); + events.push({...transformToPlayEvent(psLifecycle), playId: p.id}); } } this.emitEvent('discoveredToScrobble', { @@ -982,18 +982,18 @@ export default abstract class AbstractSource extends AbstractComponent implement try { const {lifecycle = [], ...preCompared} = await this.transformPlay(currQueuedPlay.play, TRANSFORM_HOOK.preCompare); if(lifecycle.length > 0) { - events.push({eventName: PLAY_EVENT_TYPE.transform, createdAt: dayjs(lifecycle[0].createdAt), data: lifecycle}); + events.push(transformToPlayEvent(lifecycle)); } let existing: PlayObject; // cheap check for existing const cheapExisting = await this.playRepo.checkExisting(preCompared, { notId: currQueuedPlay.id }); if(cheapExisting !== undefined) { - events.push({eventName: PLAY_EVENT_TYPE.dupeCheck, data: {match: true, closestMatchedPlay: cheapExisting.play, score: 1, breakdowns: [], createdAt: dayjs().toISOString(), reason: `Matched hash on existing Play ${cheapExisting.uid} with close temporality`}, createdAt: dayjs()}); + events.push(dupeCheckToPlayEvent({match: true, reason: `Matched hash on existing Play ${cheapExisting.uid} with close temporality`})); updatedQueueState.error = {message: `Matched hash on existing Play ${cheapExisting.uid} with close temporality`}; existing = {...cheapExisting.play, id: cheapExisting.id, uid: cheapExisting.uid}; } else { const matchRes = await this.existingDiscovered(preCompared); - events.push({eventName: PLAY_EVENT_TYPE.dupeCheck, data: matchRes, createdAt: dayjs()}); + events.push(dupeCheckToPlayEvent(matchRes)); if(matchRes.match) { existing = matchRes.closestMatchedPlay; updatedQueueState.error = {message: `Matched with Play ${existing.uid ?? existing.id}`}; @@ -1005,11 +1005,11 @@ export default abstract class AbstractSource extends AbstractComponent implement 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.'} - events.push({eventName: PLAY_EVENT_TYPE.playStateChange, data: {state, reason: 'Not added as discovered because monitoring was disabled when Play was created'}, createdAt: dayjs()}); + updatedQueueState.error = {message: 'Play was not added as discovered because monitoring was disabled when Play was created.'}; + events.push(stateChangeToPlayEvent({state, reason: 'Not added as discovered because monitoring was disabled when Play was created'})); } else { state = 'discovered'; - events.push({eventName: PLAY_EVENT_TYPE.playStateChange, data: {state}, createdAt: dayjs()}); + events.push(stateChangeToPlayEvent({state})); this.tracksDiscovered++; this.tracksDiscoveredTotal++ this.discoveredCounter.labels(this.getPrometheusLabels()).inc(); @@ -1018,7 +1018,7 @@ export default abstract class AbstractSource extends AbstractComponent implement } else { this.playRepo.updateById(existing.id, {updatedAt: dayjs()}); state = 'duped'; - events.push({eventName: PLAY_EVENT_TYPE.playStateChange, data: {state}, createdAt: dayjs()}); + events.push(stateChangeToPlayEvent({state})); currQueuedPlay.parentId = existing.id; } this.playRepo.updateById(currQueuedPlay.id, {play: preCompared, state}); @@ -1039,14 +1039,14 @@ export default abstract class AbstractSource extends AbstractComponent implement } } updatedQueueState.queueStatus = 'completed'; - events.push({eventName: PLAY_EVENT_TYPE.queueStateChange, data: queueStateToEventData({...queueState, ...updatedQueueState}), createdAt: dayjs()}); + events.push(queueStateToPlayEvent({...queueState, ...updatedQueueState})); 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'; - events.push({eventName: PLAY_EVENT_TYPE.queueStateChange, data: queueStateToEventData({...queueState, ...updatedQueueState}), createdAt: dayjs()}); + events.push(queueStateToPlayEvent({...queueState, ...updatedQueueState})); } finally { await this.queueRepo.updateById(queueState.id, updatedQueueState); await this.playEventsRepo.createMany(events.map(x => ({...x, playId: currQueuedPlay.id})));