diff --git a/src/backend/common/database/drizzle/schema/schema.ts b/src/backend/common/database/drizzle/schema/schema.ts index 19f507fa..3d3f8f67 100644 --- a/src/backend/common/database/drizzle/schema/schema.ts +++ b/src/backend/common/database/drizzle/schema/schema.ts @@ -144,7 +144,7 @@ export const queueStates = sqliteTable("play_queue_states", { queueStatus: text({enum: ['queued','completed','failed']}).notNull().default('queued'), retries: integer().notNull().default(0), error: ErrorLikeJson('error'), - context: text({mode: 'json'}).$type(), + context: text({mode: 'json'}).$type(), createdAt: DayjsTimestamp('createdAt').notNull().$defaultFn(() => dayjs()), updatedAt: DayjsTimestamp('updatedAt').notNull().$defaultFn(() => dayjs()).$onUpdate(() => dayjs()) }, (table) => [ diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index 8f907f69..1eec892e 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -1399,7 +1399,35 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i await this.updateQueueStats([DEAD_QUEUE]); } - queueScrobble = async (data: (PlayObject | PlayObject[]) | (PlaySelectWithQueueStates | PlaySelectWithQueueStates[]), context?: QueueContext) => { + public cancelQueuedPlay = async (playEntity: PlaySelectWithQueueStates) => { + const queueState = playEntity.queueStates.find(x => x.queueName === INGRESS_QUEUE); + if(queueState === undefined) { + throw new SimpleError('Play does not have an associated queued'); + } + if(queueState.queueStatus !== 'queued') { + throw new SimpleError('Play is not queued'); + } + + queueState.queueStatus = QUEUE_STATUS_FAILED; + playEntity.state = 'failed'; + const createdEvents = await this.playEventsRepo.createMany([ + {playId: playEntity.id, ...stateChangeToPlayEvent({state: playEntity.state})}, + {playId: playEntity.id, ...queueStateToPlayEvent({...queueState, context: {reason: 'Cancelled by user'}})} + ]) as PlayEventSelect[]; + await this.queueRepo.updateById(queueState.id, {queueStatus: QUEUE_STATUS_FAILED}); + await this.playRepo.updateById(playEntity.id, {state: 'failed'}); + this.emitPlayUpdate({ + ...playEntity, + events: ((playEntity as unknown as PlayWith<'events'>).events ?? []).concat(createdEvents), + } as unknown as PlayApiCommonDetailed); + if(queueState.retries === 0) { + this.emitEvent('playDequeued', { queuedScrobble: playEntity }); + } else { + this.emitEvent('deadLetterDequeued', { queuedScrobble: playEntity }); + } + } + + queueScrobble = async (data: (PlayObject | PlayObject[]) | (PlaySelectWithQueueStates | PlaySelectWithQueueStates[]), context?: QueueContext & {isRetry?: boolean}) => { const createdQueuedPlays: PlaySelect[] = []; const dataArray = Array.isArray(data) ? data : [data]; diff --git a/src/backend/server/api.ts b/src/backend/server/api.ts index 484e9a4f..ab1968b8 100644 --- a/src/backend/server/api.ts +++ b/src/backend/server/api.ts @@ -17,12 +17,13 @@ import { type SOURCE_SOT_TYPES, type SourcePlayerJson, type SourceStatusData, + queueContextSchema, } from "../../core/Atomic.ts"; import { capitalize } from "../../core/StringUtils.ts"; import type {ExpressHandler, LeveledLogData} from "../common/infrastructure/Atomic.ts"; import { getRoot } from "../ioc.ts"; import AbstractScrobbleClient from "../scrobblers/AbstractScrobbleClient.ts"; -import type AbstractSource from "../sources/AbstractSource.ts"; +import AbstractSource from "../sources/AbstractSource.ts"; import MemorySource from "../sources/MemorySource.ts"; import { parseBool } from "../utils.ts"; import { sortByNewestPlayDate } from '../../core/PlayUtils.ts'; @@ -412,7 +413,6 @@ export const setupApi = (app: Express, router: ReturnType { const { component, - query, params: { playUid } @@ -424,6 +424,62 @@ export const setupApi = (app: Express, router: ReturnType { + const { + component, + params: { + playUid + }, + body = {} + } = req; + + const play = await component.playRepo.findByUid(playUid); + if(play === undefined) { + return res.sendStatus(404); + } + + if(component instanceof AbstractSource) { + await component.queuePlay([play], {...body, isRetry: true}); + } else { + await component.queueScrobble([play], {...body, isRetry: true}); + } + return res.sendStatus(200); + }); + + router.delete('/components/:componentVal/plays/:playUid/queue', {middleware: [componentAwareMiddle]}, async (req, res, next) => { + const { + component, + params: { + playUid + } + } = req; + + const play = await component.playRepo.findByUid(playUid); + if(play === undefined) { + return res.sendStatus(404); + } + + await component.cancelQueuedPlay(play); + return res.sendStatus(200); + }); + + router.delete('/components/:componentVal/plays/:playUid/dead', {middleware: [componentAwareMiddle]}, async (req, res, next) => { + const { + component, + params: { + playUid + } + } = req; + + const play = await component.playRepo.findByUid(playUid); + if(play === undefined) { + return res.sendStatus(404); + } + + await component.removeDeadLetterScrobble(play); + return res.sendStatus(200); + }); + router.delete('/cache/:cacheType', async (req, res) => { const cache = await getRoot().items.cache(); logger.verbose(`User request cache deletion for ${req.params.cacheType}`); diff --git a/src/backend/sources/AbstractSource.ts b/src/backend/sources/AbstractSource.ts index dc3b779b..89abc1f5 100644 --- a/src/backend/sources/AbstractSource.ts +++ b/src/backend/sources/AbstractSource.ts @@ -119,7 +119,7 @@ export default abstract class AbstractSource extends AbstractComponent implement declare protected componentType: 'source'; - protected playRepo!: DrizzlePlayRepository; + public playRepo!: DrizzlePlayRepository; protected queueRepo!: DrizzleQueueRepository; protected playEventsRepo!: DrizzlePlayEventsRepository; @@ -421,7 +421,7 @@ 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[]) | (PlaySelectWithQueueStates | PlaySelectWithQueueStates[]), context?: QueueContext) => { + queuePlay = async (data: (PlayObject | PlayObject[]) | (PlaySelectWithQueueStates | PlaySelectWithQueueStates[]), context?: QueueContext & {isRetry?: boolean}) => { const createdQueuedPlays: PlaySelect[] = []; const dataArray = Array.isArray(data) ? data : [data]; @@ -1047,6 +1047,34 @@ export default abstract class AbstractSource extends AbstractComponent implement } + public cancelQueuedPlay = async (playEntity: PlaySelectWithQueueStates) => { + const queueState = playEntity.queueStates.find(x => x.queueName === INGRESS_QUEUE); + if(queueState === undefined) { + throw new SimpleError('Play does not have an associated queued'); + } + if(queueState.queueStatus !== 'queued') { + throw new SimpleError('Play is not queued'); + } + + queueState.queueStatus = QUEUE_STATUS_FAILED; + playEntity.state = 'failed'; + const createdEvents = await this.playEventsRepo.createMany([ + {playId: playEntity.id, ...stateChangeToPlayEvent({state: playEntity.state})}, + {playId: playEntity.id, ...queueStateToPlayEvent({...queueState, context: {reason: 'Cancelled by user'}})} + ]) as PlayEventSelect[]; + await this.queueRepo.updateById(queueState.id, {queueStatus: QUEUE_STATUS_FAILED}); + await this.playRepo.updateById(playEntity.id, {state: 'failed'}); + this.emitPlayUpdate({ + ...playEntity, + events: ((playEntity as unknown as PlayWith<'events'>).events ?? []).concat(createdEvents), + } as unknown as PlayApiCommonDetailed); + if(queueState.retries === 0) { + this.emitEvent('playDequeued', { queuedScrobble: playEntity }); + } else { + this.emitEvent('deadLetterDequeued', { queuedScrobble: playEntity }); + } + } + protected handlePlayProcessing = async (playEntity: PlaySelectWithQueueStates, signal?: AbortSignal) => { let res: PlayProcessingResult, err: Error; diff --git a/src/core/Atomic.ts b/src/core/Atomic.ts index 97672e2d..d0beeedf 100644 --- a/src/core/Atomic.ts +++ b/src/core/Atomic.ts @@ -657,13 +657,14 @@ export const QUEUE_STATUSES: QueueStatus[] = [QUEUE_STATUS_COMPLETED, QUEUE_STAT export const DEAD_LETTER_RETRIES_DEFAULT = 3; -export interface QueueContext { - transform?: boolean - dupeCheck?: boolean - useCache?: boolean - isRetry?: boolean - reason?: string -} +export const queueContextSchema = z.object({ + transform: z.boolean().optional(), + dupeCheck: z.boolean().optional(), + useCache: z.boolean().optional(), + reason: z.string().optional() +}); + +export type QueueContext = z.infer; /** * @see https://github.com/ts-essentials/ts-essentials/issues/339#issuecomment-4681920369 */