diff --git a/src/backend/common/AbstractComponent.ts b/src/backend/common/AbstractComponent.ts index 62706649..da33cf20 100644 --- a/src/backend/common/AbstractComponent.ts +++ b/src/backend/common/AbstractComponent.ts @@ -29,6 +29,8 @@ import { loggerNoop } from "./MaybeLogger.js"; import { objectsEqual } from "../utils/DataUtils.js"; import { RetentionOptions } from "./infrastructure/config/database.js"; import { getRetentionCompactAfterFromEnv, getRetentionDeleteAfterFromEnv, isCompactableProperty, parseRetentionOptions, parseRetentionOptionsDurations } from "./database/Database.js"; +import { DbConcrete } from "./database/drizzle/drizzleUtils.js"; +import { ComponentSelect } from "./database/drizzle/drizzleTypes.js"; export type AbstractComponentConfig = (CommonClientConfig | CommonSourceConfig) & { transformManager?: TransformerManager }; @@ -40,12 +42,15 @@ export default abstract class AbstractComponent extends AbstractInitializable { regexCache!: ReturnType; protected transformManager: TransformerManager; protected cache: MSCache; + protected db: DbConcrete; + protected dbComponent: ComponentSelect; protected retentionOpts: RetentionOptions; protected constructor(config: AbstractComponentConfig) { super(config); this.transformManager = config.transformManager ?? getRoot().items.transformerManager; this.cache = getRoot().items.cache(); + this.db = getRoot().items.db(); const cProps = config.options?.retention?.compact ?? parseArrayFromMaybeString(process.env.COMPACT_PROPERTIES, {lower: true}); if(!cProps.every(isCompactableProperty)) { throw new SimpleError(`Compactable properties must be one of 'transform' or 'input'. Given: ${cProps.join(',')}`); @@ -66,6 +71,11 @@ export default abstract class AbstractComponent extends AbstractInitializable { } } + protected async doBuildDatabase(): Promise { + super.doBuildDatabase(); + return; + } + public buildTransformRules() { this.logger.debug('Building transformer rules...'); try { diff --git a/src/backend/common/AbstractInitializable.ts b/src/backend/common/AbstractInitializable.ts index 43c68842..4bb1013a 100644 --- a/src/backend/common/AbstractInitializable.ts +++ b/src/backend/common/AbstractInitializable.ts @@ -13,6 +13,7 @@ export default abstract class AbstractInitializable { authFailure?: boolean; buildOK?: boolean | null; + databaseOK?: boolean | null; connectionOK?: boolean | null; cacheOK?: boolean | null; @@ -41,6 +42,7 @@ export default abstract class AbstractInitializable { if(this.componentLogger === undefined) { await this.buildComponentLogger(); } + await this.buildDatabase(force); await this.buildInitData(force); await this.parseCache(force); try { @@ -96,13 +98,13 @@ export default abstract class AbstractInitializable { if(!force) { return; } - this.logger.debug('Cache OK but step was forced'); + this.logger.verbose('Cache OK but step was forced'); } try { const res = await this.doParseCache(); if(res === undefined) { this.cacheOK = null; - this.logger.debug('No cache to parse.'); + this.logger.trace('No cache to parse.'); return; } if (res === true) { @@ -139,13 +141,13 @@ export default abstract class AbstractInitializable { if(!force) { return; } - this.logger.debug('Build OK but step was forced'); + this.logger.verbose('Build OK but step was forced'); } try { const res = await this.doBuildInitData(); if(res === undefined) { this.buildOK = null; - this.logger.debug('No required data to build.'); + this.logger.trace('No required data to build.'); return; } if (res === true) { @@ -172,6 +174,44 @@ export default abstract class AbstractInitializable { return; } + public async buildDatabase(force: boolean = false) { + if(this.databaseOK) { + if(!force) { + return; + } + this.logger.verbose('Database OK but step was forced'); + } + try { + const res = await this.doBuildDatabase(); + if(res === undefined) { + this.databaseOK = null; + this.logger.trace('No required database steps.'); + return; + } + if (res === true) { + this.logger.verbose('Required database init succeeded'); + } else if (typeof res === 'string') { + this.logger.verbose(`Required database init succeeded => ${res}`); + } + this.databaseOK = true; + } catch (e) { + this.databaseOK = false; + throw new BuildDataError('Required database init failed', {cause: e}); + } + } + + /** + * Run/fetch/create any database data needed for this component to operate when ready + * + * * Return undefined if not possible or not required + * * Return TRUE if database steps succeeded + * * Return string if database steps succeeded and should log result + * * Throw error on failure + * */ + protected async doBuildDatabase(): Promise { + return; + } + public async checkConnection(force: boolean = false) { if(this.connectionOK) { if(!force) { @@ -252,12 +292,14 @@ export default abstract class AbstractInitializable { public isReady() { return (this.buildOK === null || this.buildOK === true) && + (this.databaseOK === null || this.databaseOK === true) && (this.connectionOK === null || this.connectionOK === true) && !this.authGated(); } public isUsable() { return (this.buildOK === null || this.buildOK === true) && + (this.databaseOK === null || this.databaseOK === true) && (this.connectionOK === null || this.connectionOK === true); } diff --git a/src/backend/common/database/drizzle/drizzleUtils.ts b/src/backend/common/database/drizzle/drizzleUtils.ts index f8975a8b..851b7aa6 100644 --- a/src/backend/common/database/drizzle/drizzleUtils.ts +++ b/src/backend/common/database/drizzle/drizzleUtils.ts @@ -95,6 +95,22 @@ export const migrateDb = async (db: ReturnType, opts: {logger?: } } +export const migrateDbSync = (db: ReturnType, opts: {logger?: Logger, migrationsFolder?: string} = {}) => { + const { + migrationsFolder, + logger: parentLogger = loggerNoop + } = opts; + const logger = childLogger(parentLogger, 'Migrations'); + + try { + logger.info('Starting migrations...'); + migrate(db, { migrationsFolder: migrationsFolder ?? path.resolve(projectDir, 'src/backend/common/database/drizzle/migrations') }); + logger.info('Migrations complete'); + } catch (e) { + throw new Error('Failed to migrate database', { cause: e }); + } +} + export const performDbMigrationWithBackup = async (dbName: string = 'ms', opts: { logger?: Logger, workingDirectory?: string, migrationsFolder?: string } = {}) => { const dbPath = getDbPath(dbName, opts.workingDirectory); diff --git a/src/backend/common/infrastructure/config/common.ts b/src/backend/common/infrastructure/config/common.ts index d64a94ec..d3ec34f5 100644 --- a/src/backend/common/infrastructure/config/common.ts +++ b/src/backend/common/infrastructure/config/common.ts @@ -2,6 +2,13 @@ import { keyOmit } from "../Atomic.js"; export interface CommonConfig { name?: string + /** A UNIQUE identifier for this Source/Client + * + * It should be unique for the given Source/Client type. No other Source/Client of the same type should have this ID. This ID will be used to register this Source/Client in the database so that it can be identified even if you change the name of the component. + * + * If no id is given the name of this component will be used. + */ + id?: string data?: CommonData /** * Should MS use this client/source? Defaults to true diff --git a/src/backend/ioc.ts b/src/backend/ioc.ts index 2311f2fe..f805a88e 100644 --- a/src/backend/ioc.ts +++ b/src/backend/ioc.ts @@ -28,7 +28,7 @@ export interface RootOptions { cache?: CacheConfigOptions | MSCache | (() => MSCache) mbMap?: MusicBrainzSingletonMap | (() => MusicBrainzSingletonMap) transformers?: TransformerCommonConfig[] - db?: DbConcrete + db?: DbConcrete | (() => DbConcrete) } const discovered = new prom.Counter({ @@ -92,6 +92,13 @@ const createRoot = (options: RootOptions = {logger: loggerDebug}) => { maybeSingletonMb = new Map(); } + let dbFunc: () => DbConcrete; + let maybeSingletonDb: DbConcrete; + if(typeof db === 'function') { + dbFunc = db; + } else { + maybeSingletonDb = db; + } const cEmitter = new WildcardEmitter(); // do nothing, just catch @@ -151,7 +158,7 @@ const createRoot = (options: RootOptions = {logger: loggerDebug}) => { cache: () => maybeSingletonCache !== undefined ? () => maybeSingletonCache : cacheFunc, mbMap: () => maybeSingletonMb !== undefined ? () => maybeSingletonMb : mbFunc, coverArtApi, - db: db as DbConcrete + db: () => maybeSingletonDb !== undefined ? () => maybeSingletonDb : dbFunc }).add((items) => { const localUrl = generateBaseURL(baseUrl, items.port) return { diff --git a/src/backend/scrobblers/AbstractScrobbleClient.ts b/src/backend/scrobblers/AbstractScrobbleClient.ts index fdb4cca2..6062a9c8 100644 --- a/src/backend/scrobblers/AbstractScrobbleClient.ts +++ b/src/backend/scrobblers/AbstractScrobbleClient.ts @@ -69,6 +69,9 @@ import {serializeError} from 'serialize-error'; import { DEFAULT_NEW_PADDING, groupPlaysToTimeRanges } from "../utils/ListenFetchUtils.js"; import { spawn, catchAbortError, isAbortError, rethrowAbortError, delay, forever, AbortError, throwIfAborted } from 'abort-controller-x'; import { Queue, MemoryStorage } from '@platformatic/job-queue' +import { FindOne, FindWhere } from "../common/database/drizzle/drizzleTypes.js"; +import { components } from "../common/database/drizzle/schema/schema.js"; +import { generateComponentEntity } from "../common/database/drizzle/entityUtils.js"; type PlatformMappedPlays = Map; type NowPlayingQueue = Map; @@ -432,6 +435,29 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i return `Scrobbles from Cache: ${cachedQLength} Queue | ${cachedDLength} Dead Letter`; } + protected async doBuildDatabase(): Promise { + super.doBuildDatabase(); + let where: FindWhere<'components'> = { + mode: 'client', + type: this.type, + uid: this.config.id ?? this.config.name + }; + const component = await this.db.query.components.findFirst({ + where + }); + if(component !== undefined) { + this.dbComponent = component; + return; + } + + this.dbComponent = (await this.db.insert(components).values(generateComponentEntity({ + uid: this.config.id ?? this.config.name, + mode: 'client', + type: this.type, + name: this.config.name + })).returning())[0]; + } + protected async postInitialize(): Promise { const { options: { diff --git a/src/backend/sources/AbstractSource.ts b/src/backend/sources/AbstractSource.ts index 900da28a..3dc13a22 100644 --- a/src/backend/sources/AbstractSource.ts +++ b/src/backend/sources/AbstractSource.ts @@ -46,6 +46,9 @@ import prom, { Counter, Gauge } from 'prom-client'; import { normalizeStr } from '../utils/StringUtils.js'; import { spawn, catchAbortError, isAbortError, rethrowAbortError, delay, forever, AbortError, throwIfAborted } from 'abort-controller-x'; import { AbortedError, generateLoggableAbortReason } from '../common/errors/MSErrors.js'; +import { FindWhere } from '../common/database/drizzle/drizzleTypes.js'; +import { components } from '../common/database/drizzle/schema/schema.js'; +import { generateComponentEntity } from '../common/database/drizzle/entityUtils.js'; export interface RecentlyPlayedOptions { limit?: number @@ -127,6 +130,29 @@ export default abstract class AbstractSource extends AbstractComponent implement } } + protected async doBuildDatabase(): Promise { + super.doBuildDatabase(); + let where: FindWhere<'components'> = { + mode: 'source', + type: this.type, + uid: this.config.id ?? this.config.name + }; + const component = await this.db.query.components.findFirst({ + where + }); + if(component !== undefined) { + this.dbComponent = component; + return; + } + + this.dbComponent = (await this.db.insert(components).values(generateComponentEntity({ + uid: this.config.id ?? this.config.name, + mode: 'source', + type: this.type, + name: this.config.name + })).returning())[0]; + } + protected async postCache(): Promise { await super.postCache(); this.generateStaggerMappers(); diff --git a/src/backend/tests/cache/cache.test.ts b/src/backend/tests/cache/cache.test.ts index d88b5a1c..40564f20 100644 --- a/src/backend/tests/cache/cache.test.ts +++ b/src/backend/tests/cache/cache.test.ts @@ -9,7 +9,7 @@ import { generatePlays } from "../../../core/PlayTestUtils.js"; import { ListenProgressPositional, ListenProgressTS } from "../../sources/PlayerState/ListenProgress.js"; import { isPortReachableConnect } from "../../utils/NetworkUtils.js"; import { getRoot } from "../../ioc.js"; -import { transientCache } from "../utils/CacheTestUtils.js"; +import { transientCache } from "../utils/TransientTestUtils.js"; import { TestScrobbler } from "../scrobbler/TestScrobbler.js"; import { sleep } from "../../utils.js"; import {promises} from 'node:fs'; diff --git a/src/backend/tests/component/transformers.test.ts b/src/backend/tests/component/transformers.test.ts index b7766a42..26da476f 100644 --- a/src/backend/tests/component/transformers.test.ts +++ b/src/backend/tests/component/transformers.test.ts @@ -15,7 +15,7 @@ import { initMemoryCache } from "../../common/Cache.js"; import { Cacheable } from "cacheable"; import { TransformerCommonConfig } from "../../../core/Atomic.js"; import TransformerManager from "../../common/transforms/TransformerManager.js"; -import { transientCache } from "../utils/CacheTestUtils.js"; +import { transientCache } from "../utils/TransientTestUtils.js"; import dayjs from "dayjs"; import clone from "clone"; diff --git a/src/backend/tests/scrobbler/scrobblers.test.ts b/src/backend/tests/scrobbler/scrobblers.test.ts index ec711556..97879793 100644 --- a/src/backend/tests/scrobbler/scrobblers.test.ts +++ b/src/backend/tests/scrobbler/scrobblers.test.ts @@ -23,7 +23,7 @@ import { DEFAULT_CONSOLIDATE_DURATION, DEFAULT_GROUP_DURATION, groupPlaysToTimeR import { asPlay } from '../../../core/PlayMarshalUtils.js'; import { nanoid } from 'nanoid'; import { getRoot } from '../../ioc.js'; -import { transientCache } from '../utils/CacheTestUtils.js'; +import { transientCache } from '../utils/TransientTestUtils.js'; chai.use(asPromised); diff --git a/src/backend/tests/setup.ts b/src/backend/tests/setup.ts index ecc3656a..410d0a5c 100644 --- a/src/backend/tests/setup.ts +++ b/src/backend/tests/setup.ts @@ -1,6 +1,6 @@ import { loggerTest } from '@foxxmd/logging'; import { getRoot } from "../ioc.js"; -import { transientCache } from './utils/CacheTestUtils.js'; +import { transientCache, transientDb } from './utils/TransientTestUtils.js'; -const root = getRoot({cache: transientCache, logger: loggerTest}); +const root = getRoot({cache: transientCache, logger: loggerTest, db: transientDb}); root.items.cache().init(); \ No newline at end of file diff --git a/src/backend/tests/utils/CacheTestUtils.ts b/src/backend/tests/utils/CacheTestUtils.ts deleted file mode 100644 index 3b537f41..00000000 --- a/src/backend/tests/utils/CacheTestUtils.ts +++ /dev/null @@ -1,4 +0,0 @@ -import { loggerTest } from "@foxxmd/logging"; -import { MSCache } from "../../common/Cache.js"; - -export const transientCache = () => new MSCache(loggerTest, {scrobble: {provider: 'memory'}, auth: {provider: 'memory'}, metadata: {provider: 'memory'}}); \ No newline at end of file diff --git a/src/backend/tests/utils/TransientTestUtils.ts b/src/backend/tests/utils/TransientTestUtils.ts new file mode 100644 index 00000000..3f8f51da --- /dev/null +++ b/src/backend/tests/utils/TransientTestUtils.ts @@ -0,0 +1,11 @@ +import { loggerTest } from "@foxxmd/logging"; +import { MSCache } from "../../common/Cache.js"; +import { getDb, migrateDbSync } from "../../common/database/drizzle/drizzleUtils.js"; + +export const transientCache = () => new MSCache(loggerTest, { scrobble: { provider: 'memory' }, auth: { provider: 'memory' }, metadata: { provider: 'memory' } }); + +export const transientDb = () => { + const db = getDb(':memory:'); + migrateDbSync(db); + return db; +} \ No newline at end of file