diff --git a/src/backend/common/AbstractComponent.ts b/src/backend/common/AbstractComponent.ts index bab2f122..eff45e8f 100644 --- a/src/backend/common/AbstractComponent.ts +++ b/src/backend/common/AbstractComponent.ts @@ -27,6 +27,8 @@ import { diffObjects, diffObjectsConsoleOutput, patchObject } from "../../core/D import clone from "clone"; import { loggerNoop } from "./MaybeLogger.js"; import { objectsEqual } from "../utils/DataUtils.js"; +import { RetentionOptionsFull } from "./infrastructure/config/database.js"; +import { parseRetentionOptions } from "./database/Database.js"; export type AbstractComponentConfig = (CommonClientConfig | CommonSourceConfig) & { transformManager?: TransformerManager }; @@ -38,11 +40,13 @@ export default abstract class AbstractComponent extends AbstractInitializable { regexCache!: ReturnType; protected transformManager: TransformerManager; protected cache: MSCache; + protected retentionOpts: RetentionOptionsFull; protected constructor(config: AbstractComponentConfig) { super(config); this.transformManager = config.transformManager ?? getRoot().items.transformerManager; this.cache = getRoot().items.cache(); + this.retentionOpts = parseRetentionOptions(config.options.retention); } protected postCache(): Promise { diff --git a/src/backend/common/database/Database.ts b/src/backend/common/database/Database.ts index 473700b9..2da1cca0 100644 --- a/src/backend/common/database/Database.ts +++ b/src/backend/common/database/Database.ts @@ -4,12 +4,17 @@ import { promises as fs } from 'fs' import { childLogger, Logger } from '@foxxmd/logging'; import { loggerNoop } from '../MaybeLogger.js'; import { fileExists, fileOrDirectoryIsWriteable } from '../../utils/FSUtils.js'; +import { DEFAULT_RETENTION_DELETE_AFTER, RententionGranular, RetentionOptions, RetentionOptionsFull } from '../infrastructure/config/database.js'; +import { DurationValue } from '../infrastructure/Atomic.js'; +import { Duration } from 'dayjs/plugin/duration.js'; +import dayjs from 'dayjs'; +import { parseDurationFromDurationValue } from '../../utils/TimeUtils.js'; export const MEMORY_DB_NAME = ':memory:'; export const isMemoryDb = (name: string): boolean => name === MEMORY_DB_NAME; export const getDbPath = (name: string = 'ms', workingDirectory?: string): string => { - if(isMemoryDb(name)) { + if (isMemoryDb(name)) { return MEMORY_DB_NAME; } return path.resolve(workingDirectory ?? configDir, `${name}.db`); @@ -27,22 +32,79 @@ export const backupDb = async (dbName: string, opts: { logger?: Logger, workingD const dbPath = getDbPath(dbName, workingDirectory); let newDb = false; - if(dbPath !== MEMORY_DB_NAME) { - if(!fileExists(dbPath)) { + if (dbPath !== MEMORY_DB_NAME) { + if (!fileExists(dbPath)) { logger.info(`Database at ${dbPath} does not exist, will create it.`); newDb = true; } try { fileOrDirectoryIsWriteable(dbPath); } catch (e) { - throw new Error('Database path/folder is not writeable, cannot backup database', {cause: e}); + throw new Error('Database path/folder is not writeable, cannot backup database', { cause: e }); } } - if(dbPath !== MEMORY_DB_NAME && !newDb) { + if (dbPath !== MEMORY_DB_NAME && !newDb) { const backupPath = `${getDbPath(`${Date.now()}-${dbName}`, workingDirectory)}.bak`; logger.info(`Backing up database before migrating => ${backupPath}`); await fs.copyFile(dbPath, backupPath) logger.info('Backed up!'); } +} + +const parseRetentionFromEnv = (): Required> => { + const deleteAfterEnv = process.env.RETENTION_DELETE_AFTER ?? DEFAULT_RETENTION_DELETE_AFTER, + deleteCompletedEnv = process.env.RETENTION_DELETE_COMPLETED_AFTER ?? deleteAfterEnv, + deleteFailedEnv = process.env.RETENTION_DELETE_FAILED_AFTER ?? deleteAfterEnv, + deleteDupedEnv = process.env.RETENTION_DELETE_DUPED_AFTER ?? deleteAfterEnv; + + return { + completed: parseDurationFromDurationValue(deleteCompletedEnv), + failed: parseDurationFromDurationValue(deleteFailedEnv), + duped: parseDurationFromDurationValue(deleteDupedEnv) + } +} + +let retentionFromEnv: Required>; +const getRetentionFromEnv = () => { + if (retentionFromEnv === undefined) { + retentionFromEnv = parseRetentionFromEnv(); + } + return retentionFromEnv; +} + +export const parseRetentionOptions = (opts: RetentionOptions = {}): RetentionOptionsFull => { + if (typeof opts.deleteAfter === 'number' || typeof opts.deleteAfter === 'string') { + const dur = parseDurationFromDurationValue(opts.deleteAfter); + return { + deleteAfter: { + completed: dur, + duped: dur, + failed: dur + } + } + } + + const fromEnv = getRetentionFromEnv(); + if (opts.deleteAfter === undefined) { + return { + deleteAfter: fromEnv + } + } + + const { + deleteAfter: { + completed = fromEnv.completed, + failed = fromEnv.failed, + duped = fromEnv.duped + } = {} + } = opts; + + return { + deleteAfter: { + completed: dayjs.isDuration(completed) ? completed : parseDurationFromDurationValue(completed), + failed: dayjs.isDuration(failed) ? failed : parseDurationFromDurationValue(failed), + duped: dayjs.isDuration(duped) ? duped : parseDurationFromDurationValue(duped), + } + } } \ No newline at end of file diff --git a/src/backend/common/errors/MSErrors.ts b/src/backend/common/errors/MSErrors.ts index 4b15cae4..06aba95f 100644 --- a/src/backend/common/errors/MSErrors.ts +++ b/src/backend/common/errors/MSErrors.ts @@ -130,4 +130,23 @@ export const generateLoggableAbortReason = (msg: string, signal: AbortSignal): A } Error.captureStackTrace(err, generateLoggableAbortReason); return err; +} + +export class InvalidRegexError extends SimpleError { + constructor(regex: RegExp | RegExp[], val?: string, url?: string, message?: string) { + const msgParts = [ + message ?? 'Regex(es) did not match the value given.', + ]; + let regArr = Array.isArray(regex) ? regex : [regex]; + for(const r of regArr) { + msgParts.push(`Regex: ${r}`) + } + if (val !== undefined) { + msgParts.push(`Value: ${val}`); + } + if (url !== undefined) { + msgParts.push(`Sample regex: ${url}`); + } + super(msgParts.join('\r\n')); + } } \ No newline at end of file diff --git a/src/backend/common/infrastructure/Atomic.ts b/src/backend/common/infrastructure/Atomic.ts index d53a1a62..2624d784 100644 --- a/src/backend/common/infrastructure/Atomic.ts +++ b/src/backend/common/infrastructure/Atomic.ts @@ -414,4 +414,16 @@ export interface ScrobbleRangeResult { fetchedAt: Dayjs } -export const REFRESH_STALE_DEFAULT = 60; \ No newline at end of file +export const REFRESH_STALE_DEFAULT = 60; + +/** + * A duration of time + * + * May be either: + * + * * a `number` of seconds + * * a `string` containing a number and a unit of time compatible with dayjs + * + * @example [60, 3600, "1 hour", "4 days"] + */ +export type DurationValue = number | string; \ No newline at end of file diff --git a/src/backend/common/infrastructure/config/aioConfig.ts b/src/backend/common/infrastructure/config/aioConfig.ts index be1241e5..ef3dd3ee 100644 --- a/src/backend/common/infrastructure/config/aioConfig.ts +++ b/src/backend/common/infrastructure/config/aioConfig.ts @@ -5,8 +5,9 @@ import { RequestRetryOptions } from "./common.js"; import { WebhookConfig } from "./health/webhooks.js"; import { CommonSourceOptions, SourceRetryOptions } from "./source/index.js"; import { SourceAIOConfig } from "./source/sources.js"; -import { CacheConfigOptions } from "../Atomic.js"; +import { CacheConfigOptions, DurationValue } from "../Atomic.js"; import { TransformerCommonConfig } from "../../../../core/Atomic.js"; +import { RetentionOptions } from "./database.js"; export interface SourceDefaults extends CommonSourceOptions { @@ -69,6 +70,10 @@ export interface AIOConfig { cache?: CacheConfigOptions transformers?: TransformerCommonConfig[] + + database?: { + retention?: RetentionOptions + } } export interface AIOClientConfig { diff --git a/src/backend/common/infrastructure/config/client/index.ts b/src/backend/common/infrastructure/config/client/index.ts index 0a280303..a4d80cec 100644 --- a/src/backend/common/infrastructure/config/client/index.ts +++ b/src/backend/common/infrastructure/config/client/index.ts @@ -1,5 +1,7 @@ +import { DurationValue } from "../../Atomic.js"; import { PlayTransformConfig, PlayTransformOptions } from "../../Transform.js"; import { CommonConfig, CommonData, RequestRetryOptions } from "../common.js"; +import { RetentionOptions } from "../database.js"; /** * Scrobble matching (between new source track and existing client scrobbles) logging options. Used for debugging. @@ -105,6 +107,8 @@ export interface CommonClientOptions extends RequestRetryOptions, UpstreamRefres deadLetterRetries?: number playTransform?: PlayTransformOptions + + retention?: RetentionOptions } export interface CommonClientConfig extends CommonConfig { diff --git a/src/backend/common/infrastructure/config/database.ts b/src/backend/common/infrastructure/config/database.ts new file mode 100644 index 00000000..66379ae2 --- /dev/null +++ b/src/backend/common/infrastructure/config/database.ts @@ -0,0 +1,18 @@ +import { Duration } from "dayjs/plugin/duration.js"; +import { DurationValue } from "../Atomic.js"; + +export interface RententionGranular { + failed?: T + completed?: T + duped?: T +} + +export interface RetentionOptions { + deleteAfter?: T | RententionGranular +} + +export interface RetentionOptionsFull { + deleteAfter: RententionGranular +} + +export const DEFAULT_RETENTION_DELETE_AFTER = 604800; // 7 days \ No newline at end of file diff --git a/src/backend/common/infrastructure/config/source/index.ts b/src/backend/common/infrastructure/config/source/index.ts index b184bc68..aa6d1196 100644 --- a/src/backend/common/infrastructure/config/source/index.ts +++ b/src/backend/common/infrastructure/config/source/index.ts @@ -2,6 +2,8 @@ import { FileLogOptions, LogLevel } from "@foxxmd/logging"; import { PlayTransformConfig, PlayTransformOptions } from "../../Transform.js"; import { CommonConfig, CommonData, RequestRetryOptions } from "../common.js"; +import { RetentionOptions } from "../database.js"; +import { DurationValue } from "../../Atomic.js"; export interface SourceRetryOptions extends RequestRetryOptions { /** @@ -108,6 +110,8 @@ export interface CommonSourceOptions extends SourceRetryOptions { scrobbleBacklogCount?: number playTransform?: PlayTransformOptions + + retention?: RetentionOptions } export interface ManualListeningOptions { diff --git a/src/backend/scrobblers/ScrobbleClients.ts b/src/backend/scrobblers/ScrobbleClients.ts index 4ce8277f..af6d8d75 100644 --- a/src/backend/scrobblers/ScrobbleClients.ts +++ b/src/backend/scrobblers/ScrobbleClients.ts @@ -130,8 +130,11 @@ export default class ScrobbleClients { const { clients: mainConfigClientConfigs = [], clientDefaults: cd = {}, + database: { + retention + } = {}, } = aioConfig; - clientDefaults = cd; + clientDefaults = {retention, ...cd}; for (const [index, c] of mainConfigClientConfigs.entries()) { const {name = 'unnamed'} = c; if(c.type === undefined) { @@ -414,7 +417,7 @@ ${sources.join('\n')}`); if (isValidConfig !== true) { throw new Error(`Config object from ${clientConfig.source || 'unknown'} with name [${clientConfig.name || 'unnamed'}] of type [${clientConfig.type || 'unknown'}] has errors: ${isValidConfig.join(' | ')}`) }*/ - const {type, name, enable = true, source, data: d = {}} = clientConfig; + const {type, name, enable = true, source, data: d = {}, options = {}} = clientConfig; if(enable === false) { this.logger.warn({labels: [`${type} - ${name}`]}, `Client from ${source} was disabled by config`); @@ -422,41 +425,41 @@ ${sources.join('\n')}`); } // add defaults - const data = {...defaults, ...d}; + const compositeOptions = {...defaults, ...options}; let newClient; this.logger.debug({labels: [`${type} - ${name}`]}, `Constructing Client from ${source}`); switch (type) { case 'maloja': const MalojaScrobbler = (await import('./MalojaScrobbler.js')).default; - newClient = new MalojaScrobbler(name, ({...clientConfig, data} as unknown as MalojaClientConfig), notifier, this.emitter, this.logger); + newClient = new MalojaScrobbler(name, ({...clientConfig, data: d, options: compositeOptions} as unknown as MalojaClientConfig), notifier, this.emitter, this.logger); break; case 'lastfm': const LastfmScrobbler = (await import('./LastfmScrobbler.js')).default; - newClient = new LastfmScrobbler(name, {...clientConfig, data } as unknown as LastfmClientConfig, this.internalConfig, notifier, this.emitter, this.logger); + newClient = new LastfmScrobbler(name, {...clientConfig, data: d, options: compositeOptions } as unknown as LastfmClientConfig, this.internalConfig, notifier, this.emitter, this.logger); break; case 'librefm': const LibrefmScrobbler = (await import('./LibrefmScrobbler.js')).default; - newClient = new LibrefmScrobbler(name, {...clientConfig, data } as unknown as LibrefmClientConfig, this.internalConfig, notifier, this.emitter, this.logger); + newClient = new LibrefmScrobbler(name, {...clientConfig, data: d, options: compositeOptions } as unknown as LibrefmClientConfig, this.internalConfig, notifier, this.emitter, this.logger); break; case 'listenbrainz': const ListenbrainzScrobbler = (await import('./ListenbrainzScrobbler.js')).default; - newClient = new ListenbrainzScrobbler(name, {...clientConfig, data: {configDir: this.internalConfig.configDir, ...data} } as unknown as ListenBrainzClientConfig, {}, notifier, this.emitter, this.logger); + newClient = new ListenbrainzScrobbler(name, {...clientConfig, data: {configDir: this.internalConfig.configDir, ...d}, options: compositeOptions } as unknown as ListenBrainzClientConfig, {}, notifier, this.emitter, this.logger); break; case 'koito': const KoitoScrobbler = (await import('./KoitoScrobbler.js')).default; - newClient = new KoitoScrobbler(name, {...clientConfig, data: {configDir: this.internalConfig.configDir, ...data} } as unknown as KoitoClientConfig, {}, notifier, this.emitter, this.logger); + newClient = new KoitoScrobbler(name, {...clientConfig, data: {configDir: this.internalConfig.configDir, ...d}, options: compositeOptions } as unknown as KoitoClientConfig, {}, notifier, this.emitter, this.logger); break; case 'tealfm': const TealScrobbler = (await import('./TealfmScrobbler.js')).default; - newClient = new TealScrobbler(name, {...clientConfig, data: {...data}} as unknown as TealClientConfig, {}, notifier, this.emitter, this.logger); + newClient = new TealScrobbler(name, {...clientConfig, data: d, options: compositeOptions} as unknown as TealClientConfig, {}, notifier, this.emitter, this.logger); break; case 'rocksky': const RockskyScrobbler = (await import('./RockskyScrobbler.js')).default; - newClient = new RockskyScrobbler(name, {...clientConfig, data: {configDir: this.internalConfig.configDir, ...data} } as unknown as RockSkyClientConfig, {}, notifier, this.emitter, this.logger); + newClient = new RockskyScrobbler(name, {...clientConfig, data: {configDir: this.internalConfig.configDir, ...d}, options: compositeOptions } as unknown as RockSkyClientConfig, {}, notifier, this.emitter, this.logger); break; case 'discord': const DiscordScrobbler = (await import('./DiscordScrobbler.js')).default; - newClient = new DiscordScrobbler(name, {...clientConfig, data: {configDir: this.internalConfig.configDir, ...data} } as unknown as DiscordClientConfig, {}, notifier, this.emitter, this.logger); + newClient = new DiscordScrobbler(name, {...clientConfig, data: {configDir: this.internalConfig.configDir, ...d}, options: compositeOptions } as unknown as DiscordClientConfig, {}, notifier, this.emitter, this.logger); break; default: break; diff --git a/src/backend/sources/ScrobbleSources.ts b/src/backend/sources/ScrobbleSources.ts index b11bac50..db80b941 100644 --- a/src/backend/sources/ScrobbleSources.ts +++ b/src/backend/sources/ScrobbleSources.ts @@ -54,6 +54,7 @@ import { ListenBrainzData } from '../common/infrastructure/config/client/listenb import { KoitoData } from '../common/infrastructure/config/client/koito.js'; import { TealData } from '../common/infrastructure/config/client/tealfm.js'; import { RockSkyData } from '../common/infrastructure/config/client/rocksky.js'; +import { DEFAULT_RETENTION_DELETE_AFTER } from '../common/infrastructure/config/database.js'; type groupedNamedConfigs = {[key: string]: ParsedConfig[]}; @@ -230,8 +231,11 @@ export default class ScrobbleSources { const { sources: mainConfigSourcesConfigs = [], sourceDefaults: sd = {}, + database: { + retention + } = {}, } = aioConfig; - sourceDefaults = this.buildSourceDefaults(sd); + sourceDefaults = this.buildSourceDefaults({retention, ...sd}); for (const [index, c] of mainConfigSourcesConfigs.entries()) { const {name = 'unnamed'} = c; if(c.type === undefined) { diff --git a/src/backend/utils/TimeUtils.ts b/src/backend/utils/TimeUtils.ts index 15a69200..566dfe74 100644 --- a/src/backend/utils/TimeUtils.ts +++ b/src/backend/utils/TimeUtils.ts @@ -25,11 +25,15 @@ import { DEFAULT_DURATION_REPEAT_PERCENT, DEFAULT_SCROBBLE_DURATION_THRESHOLD, DEFAULT_SCROBBLE_PERCENT_THRESHOLD, + DurationValue, lowGranularitySources, ScrobbleThresholdResult, } from "../common/infrastructure/Atomic.js"; import { ScrobbleThresholds } from "../common/infrastructure/config/source/index.js"; import { formatNumber } from '../../core/DataUtils.js'; +import { InvalidRegexError, SimpleError } from "../common/errors/MSErrors.js"; +import { NamedGroup, parseRegex } from "@foxxmd/regex-buddy-core"; +import { Duration } from "dayjs/plugin/duration.js"; //dayjs.extend(isToday); @@ -399,4 +403,59 @@ export const repeatDurationPlayed = (play: PlayObject, duration: number, thresho /** Convert unix timestamp in microseconds to unix timestamp in seconds */ export const usecToUnix = (usec: number): UnixTimestamp => { return Math.floor(usec / 1000); +} + +// string must only contain ISO8601 optionally wrapped by whitespace +const ISO8601_REGEX: RegExp = /^\s*((-?)P(?=\d|T\d)(?:(\d+)Y)?(?:(\d+)M)?(?:(\d+)([DW]))?(?:T(?:(\d+)H)?(?:(\d+)M)?(?:(\d+(?:\.\d+)?)S)?)?)\s*$/; +// finds ISO8601 in any part of a string +const ISO8601_SUBSTRING_REGEX: RegExp = /((-?)P(?=\d|T\d)(?:(\d+)Y)?(?:(\d+)M)?(?:(\d+)([DW]))?(?:T(?:(\d+)H)?(?:(\d+)M)?(?:(\d+(?:\.\d+)?)S)?)?)/g; +// string must only duration optionally wrapped by whitespace +const DURATION_REGEX: RegExp = /^\s*(?