Something went wrong. Try again.
[READ-ONLY] Mirror of https://github.com/FoxxMD/multi-scrobbler. Scrobble plays from multiple sources to multiple clients docs.multi-scrobbler.app
deezer docker jellyfin koito lastfm listenbrainz maloja mopidy mpris music music-assistant plex scrobble self-hosted spotify subsonic tautulli youtube-music
Something went wrong. Try again.
27 kB · 697 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698import {loggerTest, type LogDataPretty, type Logger, type LogLevel} from "@foxxmd/logging";import type {Express} from 'express';import bsseDef from 'better-sse';import bodyParser from "body-parser";import { FixedSizeList } from 'fixed-size-list';import { PassThrough } from "node:stream";import { Transform } from "stream";import { type LogOutputConfig, queueContextSchema, logLevelStandaloneSchema,} from "../../core/Atomic.ts";import type {LeveledLogData} from "../common/infrastructure/Atomic.ts";import { getRoot } from "../ioc.ts";import AbstractScrobbleClient from "../scrobblers/AbstractScrobbleClient.ts";;import MemorySource from "../sources/MemorySource.ts";import { setupAuthRoutes } from "./auth.ts";import { setupDeezerRoutes } from "./deezerRoutes.ts";import {setupLZEndpointRoutes} from "./endpointListenbrainzRoutes.ts";import {setupLastfmEndpointRoutes} from "./endpointLastfmRoutes.ts";import { makeComponentMiddle } from "./middleware.ts";import { setupWebscrobblerRoutes } from "./webscrobblerRoutes.ts";import type ScrobbleSources from "../sources/ScrobbleSources.ts";import type ScrobbleClients from "../scrobblers/ScrobbleClients.ts";import prom from 'prom-client';import { findAuthIssue, SimpleError } from "../common/errors/MSErrors.ts";import { DrizzlePlayRepository, type QueryPlaysOpts, type QueryPlaysOptsJson } from "../common/database/drizzle/repositories/PlayRepository.ts";import AbstractHistoricalScrobbleClient from "../scrobblers/AbstractHistoricalScrobbleClient.ts";import { DrizzlePlayHistoricalRepository } from "../common/database/drizzle/repositories/PlayHistoricalRepository.ts";import {componentStateBodySchema, playStateBodySchema, type ComponentClientApiJson, type ComponentSourceApiJson} from "../../core/Api.ts";import { asDayjsHydratedObject } from "../../core/DataUtils.ts";import type {Dayjs} from "dayjs";import { asSerializablePlaySelect } from "../../core/PlayMarshalUtils.ts";import { serializeError } from "serialize-error";import { z } from 'zod';import type { createTypedRouter, TypedMiddleware } from "@minisylar/express-typed-router";import { hasMetricRepositories, registerMetrics, setMetricRepositories } from "./promMetrics.ts";import pMap from "p-map";import type { PlayWith } from "../common/database/drizzle/drizzleTypes.ts";import { stripIndents } from "common-tags";
const maxBufferSize = 300;const output: Record<number, FixedSizeList<LogDataPretty>> = {};
const jsonParser = bodyParser.json({ type: ['text/*', 'application/json'] });
const createAddToLogBuffer = (levelMap: {[p: number]: string}) => (log: LogDataPretty) => { output[log.level].add({...log, levelLabel: levelMap[log.level]});}
const getLogs = (minLevel: number, limit: number = maxBufferSize, sort: 'asc' | 'desc' = 'desc'): LeveledLogData[] => { const allLogs: LeveledLogData[][] = []; for(const level of Object.keys(output)) { if(Number.parseInt(level) >= minLevel) { allLogs.push(output[level].data); } } if(sort === 'desc') { return allLogs.flat(1).sort((a, b) => b.time - a.time).slice(0, limit); } return allLogs.flat(1).sort((a, b) => a.time - b.time).slice(0, limit);}
export interface ApiArgs { app: Express router: ReturnType<typeof createTypedRouter> scrobbleSources: ScrobbleSources scrobbleClients: ScrobbleClients}
export interface ApiOptions { logger?: Logger appLoggerStream?: PassThrough initialLogOutput?: LogDataPretty[] testMode?: boolean}
export const setupApi = (args: ApiArgs, opts: ApiOptions = {}) => { const { app, router, scrobbleSources, scrobbleClients } = args; const { logger = loggerTest, appLoggerStream = new PassThrough(), initialLogOutput = [], testMode } = opts; for(const level of Object.keys(logger.levels.labels)) { output[level] = new FixedSizeList<LeveledLogData>(maxBufferSize); }
const addToLogBuffer = createAddToLogBuffer(logger.levels.labels); for(const log of initialLogOutput) { addToLogBuffer(log); } const root = getRoot();
//let logWebLevel: LogLevel = logger.level as LogLevel || (process.env.LOG_LEVEL || 'info') as LogLevel;
const logConfig: LogOutputConfig = { level: logger.level as LogLevel || (process.env.LOG_LEVEL || 'trace') as LogLevel, sort: 'descending', limit: 50, }
let logObjectStream: Transform; try { logObjectStream = new Transform({ transform: (chunk, e, cb) => { cb(null, chunk) }, objectMode: true, allowHalfOpen: true }) } catch (e) { console.log(e); }
appLoggerStream.on('data', (log: LogDataPretty) => { addToLogBuffer(log); if(log.level >= logger.levels.values[logConfig.level]) { logObjectStream.write({message: log.line, level: log.level, levelLabel: logger.levels.labels[log.level]}); } });
const componentAwareMiddle = makeComponentMiddle(scrobbleSources, scrobbleClients);
const setLogWebSettings: TypedMiddleware = async (req, res, next) => { // @ts-expect-error logLevel not part of session const sessionLevel: LogLevel | undefined = req.session.logLevel as LogLevel | undefined; if(sessionLevel !== undefined && logConfig.level !== sessionLevel) { logConfig.level = sessionLevel; } // @ts-expect-error limit not part of session const sessionLimit: number | undefined = req.session.limit as number | undefined; if(sessionLimit !== undefined && logConfig.limit !== sessionLimit) { logConfig.limit = sessionLimit; } next(); }
router.get('/api/logs/stream', {middleware: [setLogWebSettings], tags: ['Events'], summary: 'SSE Logs'}, async (req, res) => { const session = await bsseDef.createSession(req, res); await session.stream(logObjectStream); });
router.get('/api/logs', {middleware: [setLogWebSettings], tags: ['Events'], summary: 'Get Logs'}, async (req, res) => { const slicedLog = getLogs(logger.levels.values[logConfig.level], logConfig.limit + 1, logConfig.sort === 'ascending' ? 'asc' : 'desc'); return res.json({data: slicedLog, settings: logConfig}); });
router.put('/api/logs', { middleware: [jsonParser], bodySchema: z.object({ level: logLevelStandaloneSchema.optional(), limit: z.int().positive().max(500).optional() }), tags: ['Events'], summary: 'Update Log Settings' }, async (req, res) => { logConfig.level = req.body.level as LogLevel | undefined ?? logConfig.level; logConfig.limit = req.body.limit ?? logConfig.limit; const slicedLog = getLogs(logger.levels.values[logConfig.level], logConfig.limit + 1, logConfig.sort === 'ascending' ? 'asc' : 'desc'); // @ts-expect-error logLevel not part of session req.session.logLevel = logConfig.level; // @ts-expect-error limit not part of session req.session.limit = logConfig.limit; return res.json({data: slicedLog, settings: logConfig}); });
router.get('/api/events', {querySchema: z.object({ next: z.literal('true').optional().meta({description: 'When used events are sent in new ui format'}) }), tags: ['Events'], summary: 'SSE Events'}, async (req, res) => { const { query: { next: nextQs } } = req;
const isNextapi = nextQs === 'true';
const session = await bsseDef.createSession(req, res); scrobbleSources.emitter.onAny((eventName: string, payload: any) => { if(payload !== undefined && payload.from !== undefined) { if(isNextapi) { session.push({event: eventName, ...payload}, eventName); } else { session.push({event: eventName, ...payload}, payload.from); } } }); scrobbleClients.emitter.onAny((eventName: string, payload: any) => { if(payload !== undefined && payload.from !== undefined) { if(isNextapi) { session.push({event: eventName, ...payload}, eventName); } else { session.push({event: eventName, ...payload}, payload.from); } } }); });
setupDeezerRoutes(app, logger, scrobbleSources); setupWebscrobblerRoutes(app, router, logger, scrobbleSources); setupLZEndpointRoutes(app, router, logger, scrobbleSources, scrobbleClients); setupLastfmEndpointRoutes(app, router, logger, scrobbleSources); setupAuthRoutes(app, router, logger, scrobbleSources, scrobbleClients);
router.get('/api/components', {tags: ['Source/Client'], summary: 'Get All Sources/Clients'}, async (req, res, next) => {
const sourceData = scrobbleSources.sources.filter(x => x.databaseOK).map((x) => { const { canPoll = false, polling = false, requiresAuth = false, requiresAuthInteraction = false, authed = false } = x; const base: ComponentSourceApiJson = x.getApiData(); if(!x.isReady()) { if(x.buildOK === false) { base.status = 'Initializing Data Failed'; } else if(x.connectionOK === false) { base.status = 'Communication Failed'; } else if (requiresAuth && !authed) { base.status = requiresAuthInteraction ? 'Auth Interaction Required' : 'Authentication Failed Or Not Attempted' } else { base.status = 'Not Ready'; } } else { if (canPoll) { base.status = polling ? 'Polling' : 'Idle'; } else { base.status = !x.instantiatedAt.isSame(x.lastActivityAt) ? 'Received Data' : 'Awaiting Data'; } } return base; });
const clientData = scrobbleClients.clients.filter(x => x.databaseOK).map((x) => { const { requiresAuth = false, requiresAuthInteraction = false, authed = false, scrobbling = false, } = x; const base: ComponentClientApiJson = x.getApiData();
if (!x.isReady()) { if(x.buildOK === false) { base.status = 'Initializing Data Failed'; } else if(x.connectionOK === false) { base.status = 'Communication Failed'; } else if (requiresAuth && !authed) { base.status = requiresAuthInteraction ? 'Auth Interaction Required' : 'Authentication Failed Or Not Attempted' } else { base.status = 'Not Ready'; } } else { base.status = scrobbling ? 'Running' : 'Idle'; } return base; });
return res.json([...sourceData, ...clientData]); });
router.get('/api/components/:id/players', {middleware: [componentAwareMiddle], tags: ['Source/Client'], summary: 'Get Source/Client Players'}, async (req, res, next) => { if(req.component instanceof MemorySource) { return res.json(req.component.playersToObject()); } else if(req.component instanceof AbstractScrobbleClient && req.component.nowPlayingEnabled) { return res.json(req.component.getNowPlayingPlayers()); } return res.json({}); });
router.get('/api/components/:id/players/:platformId', {middleware: [componentAwareMiddle], tags: ['Source/Client'], summary: 'Get Specific Source/Client Player'}, async (req, res, next) => { const { params: { platformId } } = req; if(req.component instanceof MemorySource) {
const player = req.component.players.get(platformId as string); if(player === undefined) { return res.status(400).json({error: `No player with platform id ${platformId} exists`}); } return res.json(player); } else if(req.component instanceof AbstractScrobbleClient && req.component.nowPlayingEnabled) { const players = req.component.getNowPlayingPlayers(); if(players[platformId as string] === undefined) { return res.status(400).json({error: `No player with platform id ${platformId} exists`}); } return res.json(players[platformId as string]); } return res.status(400).json({error: `Component does not support players`}); });
router.get('/api/components/:id', {middleware: [componentAwareMiddle], tags: ['Source/Client'], summary: 'Get Source/Client'}, async (req, res) => { const { component, } = req; return res.json(component.getApiData()); }); router.post('/api/components/:id/state', { middleware: [componentAwareMiddle, jsonParser], bodySchema: componentStateBodySchema, tags: ['Source/Client'], summary: 'Update Source/Client State'}, async (req, res) => { const { component, body: { state, reason = 'invoked by api' } } = req; switch (state) { case 'stop': try { await component.stop({ reason: new SimpleError(reason, {simple: true, shortStack: true}) }) } catch (e) { return res.status(500).json({ error: serializeError(e) }); } break; case 'start': try { await component.start({ forceInit: true }) } catch (e) { return res.status(500).json({ error: serializeError(e) }); } break; case 'restart': try { await component.restart({ forceInit: true, reason: new SimpleError(reason, {simple: true, shortStack: true}) }) } catch (e) { return res.status(500).json({ error: serializeError(e) }); } break; case 'ignore': component.monitoringActivity = component.getSystemMonitoring() === false ? undefined : false; component.emitComponentUpdate({state: component.getRunningState()}); break; case 'monitor': component.monitoringActivity = component.getSystemMonitoring() === true ? undefined : true; component.emitComponentUpdate({state: component.getRunningState()}); break; default: return res.status(400).json({ error: { message: `'state' type ${state} was not handled` } }); } return res.sendStatus(200); });
router.post('/api/components/:id/auth', { middleware: [componentAwareMiddle], tags: ['Source/Client'], summary: 'Test Source/Client Authentication' }, async (req, res, next) => { const { component, } = req; let didAuth = false; try { logger.verbose('User requested auth test'); await component.testAuth(true); component.clearErrors({predicate: x => findAuthIssue(x) !== undefined}); didAuth = true; return res.sendStatus(200); } catch (e) { component.replaceErrors(e, {predicate: x => findAuthIssue(x) !== undefined}); return res.status(500).json({error: serializeError(e)}); } finally { const data = component.getApiData(); component.emitComponentUpdate({ errors: data.errors, state: data.state, status: didAuth ? 'Authenticated successfully' : data.status }); } });
router.get('/api/components/:id/art', { middleware: [componentAwareMiddle], querySchema: z.looseObject({}), tags: ['Source/Client'], summary: 'Get External Art', description: stripIndents`Fetches and returns art from the upstream/downstream external serviceNote: this is only supported by some components.` }, async (req, res, next) => { const { component, query } = req;
if('getExternalArt' in component && typeof component.getExternalArt === 'function') { const [stream, contentType] = await component.getExternalArt(query); res.writeHead(200, {'Content-Type': contentType}); try { return stream.pipe(res); } catch (e) { logger.error(new Error(`Error occurred while trying to stream art for Component ${component.componentId}`, {cause: e})); return res.status(500).json({message: 'Error during art retrieval'}); } }
return res.status(501).json({error: {message: `Component ${component.componentId} does not support art retrieval`}}); });
router.get('/api/components/:id/plays', { middleware: [componentAwareMiddle], tags: ['Plays'], summary: 'Get Paginated Plays' }, async (req, res, next) => { const { component, query } = req;
const hydratedQuery = asDayjsHydratedObject<QueryPlaysOptsJson, QueryPlaysOpts<Dayjs>>(query); const playRes = await component.getPlaysPaginated(hydratedQuery);
// @ts-expect-error its fine playRes.data = playRes.data.map(x => asSerializablePlaySelect(x)) //PlayApiCommonDetailed // plus paginatioon return res.json(playRes); });
router.get('/api/components/:id/plays/:uid', { middleware: [componentAwareMiddle], tags: ['Plays'], summary: 'Get Play' }, async (req, res, next) => { const { component, params: { uid: playUid } } = req;
const playRes = await component.getPlayApiResponse(playUid as string); //PlayApiCommonDetailed // plus paginatioon return res.json(asSerializablePlaySelect(playRes)); });
router.delete('/api/components/:id/plays/:uid', { middleware: [componentAwareMiddle], querySchema: z.object({children: z.stringbool().optional()}), tags: ['Plays'], summary: 'Delete Play' }, async (req, res, next) => { const { component, query: { children }, params: { uid: playUid } } = req;
const play = await component.playRepo.findByUidWith<'children'>(playUid, ['children']); if(play === undefined) { return res.sendStatus(404); }
await component.deletePlay(play);
return res.sendStatus(200); });
router.post('/api/components/:id/plays/queue', { middleware: [componentAwareMiddle,jsonParser], bodySchema: z.object({ context: queueContextSchema.optional(), filters: z.looseObject({}) }), tags: ['Plays'], summary: 'Requeue Bulk Plays' }, async (req, res, next) => { const { component, body } = req;
const hydratedQuery = asDayjsHydratedObject<QueryPlaysOptsJson, QueryPlaysOpts<Dayjs>>({...body.filters, with: ['queues']}); res.sendStatus(200);
const queueFunc = async (p: PlayWith<'queueStates'>) => await component.queuePlay([p], {...body.context, isRetry: true, reason: 'User requested reprocessing'});
const currentFilters = hydratedQuery; let more = true; while(more) { const res = await component.getPlaysPaginatedInternal(currentFilters); pMap(res.data, async (x) => await queueFunc(x), {concurrency: 5}); more = res.data.length === res.meta.limit; if(more) { currentFilters.offset += res.meta.limit } } });
router.post('/api/components/:id/plays/:uid/queue', { middleware: [componentAwareMiddle,jsonParser], bodySchema: queueContextSchema.optional(), tags: ['Plays'], summary: 'Requeue a Play' }, async (req, res, next) => { const { component, params: { uid: playUid }, body = {} } = req;
const play = await component.playRepo.findByUid(playUid); if(play === undefined) { return res.sendStatus(404); }
await component.queuePlay([play], {...body, isRetry: true, reason: 'User requested reprocessing'}); return res.sendStatus(200); });
router.delete('/api/components/:id/plays/:uid/queue', { middleware: [componentAwareMiddle], tags: ['Plays'], summary: 'Dequeue a Play' }, async (req, res, next) => { const { component, params: { uid: playUid } } = req;
const play = await component.playRepo.findByUidWith<'queueStates' | 'events'>(playUid, ['queues','events']); if(play === undefined) { return res.sendStatus(404); }
await component.cancelQueuedPlay(play); return res.sendStatus(200); });
router.post('/api/components/:id/plays/:uid/state', { middleware: [componentAwareMiddle, jsonParser], bodySchema: playStateBodySchema, tags: ['Plays'], summary: 'Mark Play as Done' }, async (req, res, next) => { const { component, params: { uid: playUid }, body: { state } } = req;
const play = await component.playRepo.findByUidWith<'queueStates' | 'events'>(playUid, ['queues','events']); if(play === undefined) { return res.sendStatus(404); }
switch(state) { case 'discarded': await component.markPlayDiscarded(play); break; default: return res.status(400).json({error: {message: `Play state '${state}' is not supported.`}}); }
return res.sendStatus(200); });
router.delete('/api/cache/:cacheType', { tags: ['Cache'], summary: 'Delete Cache By Type' }, async (req, res) => { const cache = await getRoot().items.cache(); logger.verbose(`User request cache deletion for ${req.params.cacheType}`); switch(req.params.cacheType) { case 'external-api': await cache.cacheApi.clear(); break; case 'transforms': await cache.cacheTransform.clear(); break; default: return res.sendStatus(404); } logger.verbose('Cache cleared!'); return res.sendStatus(204); });
router.post('/api/components/:id/plays/historical', { middleware: [componentAwareMiddle,jsonParser], tags: ['Plays'], summary: 'Hydrate Historical Plays', bodySchema: z.object({type: z.enum(['full','recent']).optional().meta({ description: `Use \`full\` to initiate a complete sync of all history or \`recent\` to sync the last ~100 plays.` })}).optional(), description: 'If the Source/Client supports Historical Play capabilities, this route requests a manual hydration of historical Plays' }, async (req, res, next) => { const { component, body: { type: syncType = 'recent' } } = req;
if(component instanceof AbstractHistoricalScrobbleClient) { component.logger.info('User requested historical play hydration'); if(syncType === 'full') { component.hydrateHistoricalScrobbles(); } else { component.syncRecentHistoricalScrobbles() .then(() => null) .catch((e) => component.logger.error(new SimpleError('Sync attempt failed', {cause: e}))); } res.status(200).send('OK'); } else { component.logger.warn('This client does not have historical play capabilities'); return res.status(409).json({error: 'This client does not have historical play capabilities'}); } });
router.get('/health', {hidden: true}, async (req, res) => res.redirect(307, `/api/${req.url.slice(1)}`)); router.get('/api/health', {querySchema: z.object({ type: z.string().optional().meta({description: 'Only report status for components of this type'}), name: z.string().optional().meta({description: 'Only report status for the component with this name'}) }), tags: ['App Meta']}, async (req, res) => { const { type, name } = req.query;
const [sourcesReady, sourceMessages] = await scrobbleSources.getStatusSummary(type, name); const [clientsReady, clientMessages] = await scrobbleClients.getStatusSummary(type, name);
return res.status((clientsReady && sourcesReady) ? 200 : 500).json({messages: sourceMessages.concat(clientMessages)}); });
if(testMode !== true) { registerMetrics(scrobbleSources, scrobbleClients); if(process.env.PROMETHEUS_FULL === 'true') { prom.collectDefaultMetrics(); } }
router.get('/api/metrics', {tags: ['App Meta']}, async (req, res) => {
if(!hasMetricRepositories()) { const db = await getRoot().items.db(); setMetricRepositories(new DrizzlePlayRepository(db),new DrizzlePlayHistoricalRepository(db)) }
const metricsString = await prom.register.metrics(); return res .status(200) .set('Content-Type', 'text/plain') .send(metricsString);
});
router.get('/api/version', {tags: ['App Meta']}, async (req, res) => { return res.json({version: root.get('version')}); });
router.use('/api/docs', router.docs({ title: "Multi-Scrobbler API", version: "0.1.0", description: "Public API docs", }))
router.all('/api/*path', {hidden: true}, async (req, res) => { const remote = req.connection.remoteAddress; const proxyRemote = req.headers["x-forwarded-for"]; const ua = req.headers["user-agent"]; logger.debug(`Server received ${req.method} request from ${remote}${proxyRemote !== undefined ? ` (${proxyRemote})` : ''}${ua !== undefined ? ` (UA: ${ua})` : ''} to unknown route: ${req.originalUrl}`); return res.sendStatus(404); });}