diff --git a/.devcontainer/devcontainer.json b/.devcontainer/devcontainer.json index 3dc2ea9f..bc52fd1e 100644 --- a/.devcontainer/devcontainer.json +++ b/.devcontainer/devcontainer.json @@ -17,7 +17,7 @@ // "features": {}, // Use 'forwardPorts' to make a list of ports inside the container available locally. - "forwardPorts": [9078,3000], + "forwardPorts": [9078,3000,5433], // Use 'postCreateCommand' to run commands after the container is created. "postCreateCommand": ".devcontainer/postCreate.sh", @@ -42,6 +42,9 @@ }, "9078": { "label": "App" + }, + "5433": { + "label": "PGLite Server" } } diff --git a/package-lock.json b/package-lock.json index 2d9d7467..e0cf2566 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,6 +16,7 @@ "@atproto/oauth-client-node": "^0.3.10", "@donedeal0/superdiff": "^1.1.1", "@electric-sql/pglite": "^0.4.5", + "@electric-sql/pglite-socket": "^0.1.5", "@ewanc26/tid": "^1.0.2", "@foxxmd/chromecast-client": "^1.0.4", "@foxxmd/get-version": "^0.0.3", @@ -1091,6 +1092,18 @@ "dev": true, "license": "Apache-2.0" }, + "node_modules/@electric-sql/pglite-socket": { + "version": "0.1.5", + "resolved": "https://registry.npmjs.org/@electric-sql/pglite-socket/-/pglite-socket-0.1.5.tgz", + "integrity": "sha512-/RAye+3EPKfO9nY4tljzxXmkT7yIpFDm0L3F+c28b+Z6uxPOjy/Zz/QEHYHXcrfuUC88/a9S72EO0+3E0j97wQ==", + "license": "Apache-2.0", + "bin": { + "pglite-server": "dist/scripts/server.js" + }, + "peerDependencies": { + "@electric-sql/pglite": "0.4.5" + } + }, "node_modules/@emotion/babel-plugin": { "version": "11.13.5", "dev": true, diff --git a/package.json b/package.json index 412eb318..9b90532e 100644 --- a/package.json +++ b/package.json @@ -16,6 +16,7 @@ "build:backend": "tsc -p src/backend && npm run -s schema:app", "build": "npm run -s build:backend && npm run -s build:frontend && npm run -s docs:build", "build:parallel": "concurrently --kill-others-on-fail --names backend,frontend,docs \"npm run -s build:backend\" \"npm run -s build:frontend\" \"npm run docs:build\"", + "db:start": "pglite-server --db=./config/msDb --port=5433 --host=0.0.0.0 -m 10", "docs:install": "cd docsite && npm install --no-audit", "docs:start": "cd docsite && npm start", "docs:build": "npm run -s schema:docs && cd docsite && npm run build", @@ -54,6 +55,7 @@ "@atproto/oauth-client-node": "^0.3.10", "@donedeal0/superdiff": "^1.1.1", "@electric-sql/pglite": "^0.4.5", + "@electric-sql/pglite-socket": "^0.1.5", "@ewanc26/tid": "^1.0.2", "@foxxmd/chromecast-client": "^1.0.4", "@foxxmd/get-version": "^0.0.3", diff --git a/src/backend/common/infrastructure/Atomic.ts b/src/backend/common/infrastructure/Atomic.ts index eb607941..a4fe2898 100644 --- a/src/backend/common/infrastructure/Atomic.ts +++ b/src/backend/common/infrastructure/Atomic.ts @@ -438,4 +438,6 @@ export const REFRESH_STALE_DEFAULT = 60; * * @example [60, 3600, "1 hour", "4 days"] */ -export type DurationValue = number | string; \ No newline at end of file +export type DurationValue = number | string; + +export type DbExternalMode = 'none' | 'live' | 'standalone'; \ No newline at end of file diff --git a/src/backend/index.ts b/src/backend/index.ts index 32e4f46c..a023f87a 100644 --- a/src/backend/index.ts +++ b/src/backend/index.ts @@ -16,16 +16,18 @@ import { appLogger, initLogger as getInitLogger } from "./common/logging.js"; import { getRoot } from "./ioc.js"; import { parseVersion } from "./version.js"; import { initServer } from "./server/index.js"; -import { isDebugMode, parseBool, retry, sleep } from "./utils.js"; +import { isDebugMode, parseBool, parseBoolStrict, retry, sleep } from "./utils.js"; import { readJson } from './utils/DataUtils.js'; import ScrobbleClients from './scrobblers/ScrobbleClients.js'; import ScrobbleSources from './sources/ScrobbleSources.js'; import { Notifiers } from './notifier/Notifiers.js'; -import { getMigratedDb } from './common/database/drizzle/drizzleUtils.js'; +import { DbConcrete, getMigratedDb } from './common/database/drizzle/drizzleUtils.js'; import { getDbPath } from './common/database/Database.js'; import { createRetentionCleanupTask } from './tasks/retentionCleanup.js'; import { parseUserConfig } from './common/Cache.js'; import { nonEmptyStringOrDefault } from '../core/StringUtils.js'; +import { DbExternalMode } from './common/infrastructure/Atomic.js'; +import { PGLiteSocketServer } from '@electric-sql/pglite-socket'; dayjs.extend(utc) dayjs.extend(isBetween); @@ -51,13 +53,48 @@ output = output.slice(0, 301); let logger: FoxLogger; -process.on('uncaughtExceptionMonitor', (err, origin) => { +let server: PGLiteSocketServer; +let db: DbConcrete; +let dbConnectionsClosed = false; + +process.on('uncaughtExceptionMonitor', async (err, origin) => { const appError = new Error(`Uncaught exception is crashing the app! :( Type: ${origin}`, {cause: err}); if(logger !== undefined) { logger.error(appError) } else { initLogger.error(appError); } + if(!dbConnectionsClosed) { + const parts = []; + if(server !== undefined) { + await server.stop(); + parts.push('PGLite Socket Server'); + } + if(db !== undefined && !db.$client.closed) { + await db.$client.close(); + parts.push('Database'); + } + if(parts.length > 0 && logger !== undefined) { + logger.info(`Closed ${parts.join(' and ')}`); + } + } +}); +process.on('SIGINT', async () => { + if(!dbConnectionsClosed) { + const parts = []; + if(server !== undefined) { + await server.stop(); + parts.push('PGLite Socket Server'); + } + if(db !== undefined && !db.$client.closed) { + await db.$client.close(); + parts.push('Database'); + } + if(parts.length > 0 && logger !== undefined) { + logger.info(`Closed ${parts.join(' and ')}`); + } + } + process.exit(0); }) const configDir = process.env.CONFIG_DIR || path.resolve(projectDir, `./config`); @@ -97,125 +134,158 @@ const configDir = process.env.CONFIG_DIR || path.resolve(projectDir, `./config`) const [aLogger, appLoggerStream] = await appLogger(logging) logger = childLogger(aLogger, 'App'); + + const dbModeVal: string = nonEmptyStringOrDefault(process.env.DB_MODE, undefined); + let dbMode: DbExternalMode; + if(dbModeVal !== undefined) { + if(['none','live','standalone'].includes(dbModeVal.toLocaleLowerCase())) { + dbMode = dbModeVal as typeof dbMode; + } else { + throw new Error(`DB_MODE env must be one of 'none' 'live' 'standalone', found ${dbModeVal}`); + } + } else { + dbMode = 'none'; + } + logger.info(`DB External Mode: ${dbMode}`); + const dbPath = getDbPath('msDb'); logger.info(`Using database at ${getDbPath('msDb')}`); - const [db, isNew] = await getMigratedDb(dbPath, {logger: childLogger(logger, 'DB')}); + const [migratedDb, isNew] = await getMigratedDb(dbPath, {logger: childLogger(logger, 'DB')}); + db = migratedDb; - const root = getRoot({ - ...config, - cache: parseUserConfig(cache, logger), - logger, - loggingConfig: logging, - loggerStream: appLoggerStream, - db - }); + - const internalConfigOptional = { - localUrl: root.get('localUrl'), - configDir: root.get('configDir'), - version: root.get('version') - }; + if(['live','standalone'].includes(dbMode)) { + server = new PGLiteSocketServer({ + host: '0.0.0.0', + port: 5433, + db: db.$client, + maxConnections: 10, + debug: parseBoolStrict(nonEmptyStringOrDefault(process.env.DB_DEBUG, false)) + }); + await server.start(); + logger.info('Started PGLite Socket Server'); + } - const scrobbleClients = new ScrobbleClients(root.get('clientEmitter'), root.get('sourceEmitter'), internalConfigOptional, root.get('logger')); - const scrobbleSources = new ScrobbleSources(root.get('sourceEmitter'), internalConfigOptional, root.get('logger')); + if (dbMode === 'standalone') { + logger.info('MS App startup stopped early due to Standalone DB Mode.'); + } else { - await root.items.cache().init(true); + const root = getRoot({ + ...config, + cache: parseUserConfig(cache, logger), + logger, + loggingConfig: logging, + loggerStream: appLoggerStream, + db + }); - initServer(logger, appLoggerStream, output, scrobbleSources, scrobbleClients); + const internalConfigOptional = { + localUrl: root.get('localUrl'), + configDir: root.get('configDir'), + version: root.get('version') + }; - if(process.env.IS_LOCAL === 'true') { - logger.info('multi-scrobbler can be run as a background service! See: https://docs.multi-scrobbler.app/installation/service'); - } + const scrobbleClients = new ScrobbleClients(root.get('clientEmitter'), root.get('sourceEmitter'), internalConfigOptional, root.get('logger')); + const scrobbleSources = new ScrobbleSources(root.get('sourceEmitter'), internalConfigOptional, root.get('logger')); - if(appConfigFail !== undefined) { - logger.warn('App config file exists but could not be parsed!'); - logger.warn(appConfigFail); - } + await root.items.cache().init(true); - const notifiers = new Notifiers(root.get('notifierEmitter'), root.get('clientEmitter'), root.get('sourceEmitter'), root.get('logger')); //root.get('notifiers'); - await notifiers.buildWebhooks(webhooks); - - await root.items.transformerManager.registerFromEnv(); - await root.items.transformerManager.registeryDefaults(); - await root.items.transformerManager.initTransformers(); - - /* - * setup clients - * */ - await scrobbleClients.buildClientsFromConfig(notifiers); - /* - * setup sources - * */ - await scrobbleSources.buildSourcesFromConfig([]); - - // check ambiguous client/source types like this for now - const lastfmSources = scrobbleSources.getByType('lastfm'); - const lastfmScrobbles = scrobbleClients.getByType('lastfm'); - - const scrobblerNames = lastfmScrobbles.map(x => x.name); - const nameColl = lastfmSources.filter(x => scrobblerNames.includes(x.name)); - if(nameColl.length > 0) { - logger.warn(`Last.FM source and clients have same names [${nameColl.map(x => x.name).join(',')}] -- this may cause issues`); - } - const clientInitOptions = {deadDelay: nonEmptyStringOrDefault(process.env.DEBUG_DEAD_DELAY, undefined) !== undefined ? Number.parseInt(process.env.DEBUG_DEAD_DELAY) : undefined}; for(const c of scrobbleClients.clients) { - c.initTasks(clientInitOptions); - const res = await Promise.race([ - sleep(2200), - (async () => { - while(!c.isReady()) { - await sleep(400) - } - return true; - })() - ]); - if(res === undefined) { - logger.debug(`Not waiting for Client ${c.name} to finish init, moving on to the next Client...`); + initServer(logger, appLoggerStream, output, scrobbleSources, scrobbleClients); + + if (process.env.IS_LOCAL === 'true') { + logger.info('multi-scrobbler can be run as a background service! See: https://docs.multi-scrobbler.app/installation/service'); } - } - for(const c of scrobbleSources.sources) { - c.initTasks(); - const res = await Promise.race([ - sleep(2200), - (async () => { - while(!c.isReady()) { - await sleep(400) - } - return true; - })() - ]); - if(res === undefined) { - logger.debug(`Not waiting for Source ${c.name} to finish init, moving on to the next Source...`); + if (appConfigFail !== undefined) { + logger.warn('App config file exists but could not be parsed!'); + logger.warn(appConfigFail); } - } - let runRetentionNow = parseBool(process.env.RETENTION_IMMEDIATE, false); - - const retentionTask = createRetentionCleanupTask(scrobbleSources, scrobbleClients, logger); - let retentionJobAdded = false; - const addJob = () => { - retentionJobAdded = true; - scheduler.addSimpleIntervalJob(new SimpleIntervalJob({ - minutes: 60, - runImmediately: runRetentionNow - }, retentionTask, {id: 'retention', preventOverrun: true})); - logger.debug('Added Retention Cleanup task to scheduler'); - }; - logger.debug('Added Client Heartbeat task to scheduler'); - - if(runRetentionNow === false || (scrobbleClients.clients.every(x => x.isReady()) && scrobbleSources.sources.every(x => x.isReady()))) { - addJob(); - } + const notifiers = new Notifiers(root.get('notifierEmitter'), root.get('clientEmitter'), root.get('sourceEmitter'), root.get('logger')); //root.get('notifiers'); + await notifiers.buildWebhooks(webhooks); + + await root.items.transformerManager.registerFromEnv(); + await root.items.transformerManager.registeryDefaults(); + await root.items.transformerManager.initTransformers(); + + /* + * setup clients + * */ + await scrobbleClients.buildClientsFromConfig(notifiers); + /* + * setup sources + * */ + await scrobbleSources.buildSourcesFromConfig([]); + + // check ambiguous client/source types like this for now + const lastfmSources = scrobbleSources.getByType('lastfm'); + const lastfmScrobbles = scrobbleClients.getByType('lastfm'); + + const scrobblerNames = lastfmScrobbles.map(x => x.name); + const nameColl = lastfmSources.filter(x => scrobblerNames.includes(x.name)); + if (nameColl.length > 0) { + logger.warn(`Last.FM source and clients have same names [${nameColl.map(x => x.name).join(',')}] -- this may cause issues`); + } + const clientInitOptions = { deadDelay: nonEmptyStringOrDefault(process.env.DEBUG_DEAD_DELAY, undefined) !== undefined ? Number.parseInt(process.env.DEBUG_DEAD_DELAY) : undefined }; for (const c of scrobbleClients.clients) { + c.initTasks(clientInitOptions); + const res = await Promise.race([ + sleep(2200), + (async () => { + while (!c.isReady()) { + await sleep(400) + } + return true; + })() + ]); + if (res === undefined) { + logger.debug(`Not waiting for Client ${c.name} to finish init, moving on to the next Client...`); + } + } + + for (const c of scrobbleSources.sources) { + c.initTasks(); + const res = await Promise.race([ + sleep(2200), + (async () => { + while (!c.isReady()) { + await sleep(400) + } + return true; + })() + ]); + if (res === undefined) { + logger.debug(`Not waiting for Source ${c.name} to finish init, moving on to the next Source...`); + } + } - logger.info('Scheduler started.'); + let runRetentionNow = parseBool(process.env.RETENTION_IMMEDIATE, false); + + const retentionTask = createRetentionCleanupTask(scrobbleSources, scrobbleClients, logger); + let retentionJobAdded = false; + const addJob = () => { + retentionJobAdded = true; + scheduler.addSimpleIntervalJob(new SimpleIntervalJob({ + minutes: 60, + runImmediately: runRetentionNow + }, retentionTask, { id: 'retention', preventOverrun: true })); + logger.debug('Added Retention Cleanup task to scheduler'); + }; + logger.debug('Added Client Heartbeat task to scheduler'); + + if (runRetentionNow === false || (scrobbleClients.clients.every(x => x.isReady()) && scrobbleSources.sources.every(x => x.isReady()))) { + addJob(); + } - if(runRetentionNow === true && !retentionJobAdded) { - logger.info('Detected that Retention Cleanup should run immediately but all sources/clients have not started yet! Delaying retention cleanup by 1 minute to allow all sources/clients to finish starting.'); - await sleep(60 * 1000); - addJob(); - } + logger.info('Scheduler started.'); + if (runRetentionNow === true && !retentionJobAdded) { + logger.info('Detected that Retention Cleanup should run immediately but all sources/clients have not started yet! Delaying retention cleanup by 1 minute to allow all sources/clients to finish starting.'); + await sleep(60 * 1000); + addJob(); + } + } } catch (e) { const appError = new Error('Exited with uncaught error', {cause: e}); @@ -224,6 +294,20 @@ const configDir = process.env.CONFIG_DIR || path.resolve(projectDir, `./config`) } else { initLogger.error(appError); } + if(!dbConnectionsClosed) { + const parts = []; + if(server !== undefined) { + await server.stop(); + parts.push('PGLite Socket Server'); + } + if(db !== undefined && !db.$client.closed) { + await db.$client.close(); + parts.push('Database'); + } + if(parts.length > 0 && logger !== undefined) { + logger.info(`Closed ${parts.join(' and ')}`); + } + } process.exit(1); } }()); diff --git a/src/core/StringUtils.ts b/src/core/StringUtils.ts index f30f5687..f78e63e0 100644 --- a/src/core/StringUtils.ts +++ b/src/core/StringUtils.ts @@ -230,7 +230,7 @@ export const splitByFirstRegexFound = (str: any, onNotAStringVal: T, delimsRe /** * Returns value if it is a non-empty string or returns default value * */ -export const nonEmptyStringOrDefault = (str: any, defaultVal: T = undefined): string | T => { +export const nonEmptyStringOrDefault = (str: any, defaultVal: T = undefined): string | T => { if (str === undefined || str === null || typeof str !== 'string' || str.trim() === '') { return defaultVal; }