Compare commits

...
33 changed files with 1213 additions and 639 deletions
+4 -1
View File
@@ -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"
}
}
+2
View File
@@ -122,6 +122,7 @@ config/*.json
config/mscache
config/yti-*
config/*.cache
config/msDb
*.txt
!robots.txt
.idea/
@@ -138,6 +139,7 @@ docsite/static/schemas/*.json
*.bak
*.bak.used
*.p8
.flatpak-builder
flatpak/generated-sources.json
+2 -1
View File
@@ -6,5 +6,6 @@
"file": [
"./src/backend/tests/setup.ts"
],
"exit": true
"exit": true,
"timeout": 2500
}
+2 -1
View File
@@ -13,5 +13,6 @@
"files.associations": {
"*.css": "tailwindcss"
},
"tailwindCSS.experimental.configFile": "src/client/index.css"
"tailwindCSS.experimental.configFile": "src/client/index.css",
"pgliteExplorer.databasePaths": []
}
+2 -2
View File
@@ -6,8 +6,8 @@ import * as path from 'path';
export default defineConfig({
schema: path.resolve(projectDir, 'src/backend/common/database/drizzle/schema'),
out: path.resolve(projectDir, 'src/backend/common/database/drizzle/migrations'),
dialect: 'sqlite',
dialect: 'postgresql',
dbCredentials: {
url: path.resolve(configDir, process.env.DB_FILE_NAME! ?? 'ms.db'),
url: path.resolve(configDir, 'msDb'),
},
});
+36 -8
View File
@@ -15,6 +15,8 @@
"@atproto/api": "^0.18.0",
"@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",
@@ -47,7 +49,7 @@
"dbus-ts": "^0.0.7",
"discord.js": "^14.26.0",
"dotenv": "^10.0.0",
"drizzle-orm": "^1.0.0-rc.2-8feace6",
"drizzle-orm": "^1.0.0-rc.2-c5a84d1",
"express": "^5.2.1",
"express-session": "^1.19.0",
"fast-equals": "^6.0.0",
@@ -104,6 +106,7 @@
"@chromatic-com/storybook": "^5.0.1",
"@curvenote/ansi-to-react": "^7.0.0",
"@dbus-types/notifications": "^0.0.5",
"@electric-sql/pglite-prepopulatedfs": "^0.0.3",
"@emotion/react": "^11.14.0",
"@eslint/js": "^8.56.0",
"@faker-js/faker": "^9.0.1",
@@ -146,7 +149,7 @@
"chai": "^4.3.6",
"chai-as-promised": "^8.0.2",
"clsx": "^2.1.1",
"drizzle-kit": "^1.0.0-rc.1-929a083",
"drizzle-kit": "^1.0.0-rc.2-c5a84d1",
"eslint": "^8.56.0",
"eslint-plugin-prefer-arrow-functions": "^3.2.4",
"eslint-plugin-storybook": "10.1.11",
@@ -1076,6 +1079,31 @@
"url": "https://github.com/sponsors/dword-design"
}
},
"node_modules/@electric-sql/pglite": {
"version": "0.4.5",
"resolved": "https://registry.npmjs.org/@electric-sql/pglite/-/pglite-0.4.5.tgz",
"integrity": "sha512-aGG2zGEyZzGWKy8P+9ZoNUV0jxt1+hgbeTf+bVAYyxVZZLXg3/9aFlfLxb08AYZVAfAkQlQIysmWjhc5hwDG8g==",
"license": "Apache-2.0"
},
"node_modules/@electric-sql/pglite-prepopulatedfs": {
"version": "0.0.3",
"resolved": "https://registry.npmjs.org/@electric-sql/pglite-prepopulatedfs/-/pglite-prepopulatedfs-0.0.3.tgz",
"integrity": "sha512-3MNFt+gR0P22foWi55j/HZ6DvQ82DEVIvmoKVYCdoG/gezMikGR794tO07/15CV4RcR3PnfCPoQjApPfXfD01w==",
"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,
@@ -7262,9 +7290,9 @@
}
},
"node_modules/drizzle-kit": {
"version": "1.0.0-rc.1-929a083",
"resolved": "https://registry.npmjs.org/drizzle-kit/-/drizzle-kit-1.0.0-rc.1-929a083.tgz",
"integrity": "sha512-Bo7T8db9V+nYC5wQjNsHgGKyxZWJkMUVUKVqveEtKxtmV3opCoG1FoHEcRTrAAz+CBzCkyM82VG2aObT229oxw==",
"version": "1.0.0-rc.2-c5a84d1",
"resolved": "https://registry.npmjs.org/drizzle-kit/-/drizzle-kit-1.0.0-rc.2-c5a84d1.tgz",
"integrity": "sha512-TjXBd6/Jo8aqC2uBbbgPoIVE5Qr9s0KubVtdj3Ro7bF2lz1TkGOhnDIzsOgMJPMSPn9uyhNG6WBKdGz+XzcL9g==",
"dev": true,
"license": "MIT",
"dependencies": {
@@ -7763,9 +7791,9 @@
}
},
"node_modules/drizzle-orm": {
"version": "1.0.0-rc.2-8feace6",
"resolved": "https://registry.npmjs.org/drizzle-orm/-/drizzle-orm-1.0.0-rc.2-8feace6.tgz",
"integrity": "sha512-5yLmJTezMNI1gUEeGf0TJAr05nK1BSs82w/L3wJNHFlnQXSFbMdxQXGcDAivGGQEuiBwZTl7gMgDNsrXXQnZIg==",
"version": "1.0.0-rc.2-c5a84d1",
"resolved": "https://registry.npmjs.org/drizzle-orm/-/drizzle-orm-1.0.0-rc.2-c5a84d1.tgz",
"integrity": "sha512-2nw0eVFNcLRU7mhs/f/UodwvRIP78ZzZXMrl/r6LpTpcJ42ZzQZZTH0ObPCR722SaEzn35kT4gvuoiam/OCK9g==",
"license": "Apache-2.0",
"peerDependencies": {
"@aws-sdk/client-rds-data": ">=3",
+6 -2
View File
@@ -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",
@@ -53,6 +54,8 @@
"@atproto/api": "^0.18.0",
"@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",
@@ -85,7 +88,7 @@
"dbus-ts": "^0.0.7",
"discord.js": "^14.26.0",
"dotenv": "^10.0.0",
"drizzle-orm": "^1.0.0-rc.2-8feace6",
"drizzle-orm": "^1.0.0-rc.2-c5a84d1",
"express": "^5.2.1",
"express-session": "^1.19.0",
"fast-equals": "^6.0.0",
@@ -142,6 +145,7 @@
"@chromatic-com/storybook": "^5.0.1",
"@curvenote/ansi-to-react": "^7.0.0",
"@dbus-types/notifications": "^0.0.5",
"@electric-sql/pglite-prepopulatedfs": "^0.0.3",
"@emotion/react": "^11.14.0",
"@eslint/js": "^8.56.0",
"@faker-js/faker": "^9.0.1",
@@ -184,7 +188,7 @@
"chai": "^4.3.6",
"chai-as-promised": "^8.0.2",
"clsx": "^2.1.1",
"drizzle-kit": "^1.0.0-rc.1-929a083",
"drizzle-kit": "^1.0.0-rc.2-c5a84d1",
"eslint": "^8.56.0",
"eslint-plugin-prefer-arrow-functions": "^3.2.4",
"eslint-plugin-storybook": "10.1.11",
+4 -4
View File
@@ -47,8 +47,8 @@ export default abstract class AbstractComponent extends AbstractInitializable {
protected transformManager: TransformerManager;
protected cache: MSCache;
protected db: DbConcrete;
protected componentRepo: DrizzleComponentRepository;
protected dbComponent: ComponentSelect;
protected componentRepo!: DrizzleComponentRepository;
protected dbComponent!: ComponentSelect;
protected retentionOpts: RetentionOptions;
protected componentType: 'source' | 'client';
@@ -58,8 +58,6 @@ export default abstract class AbstractComponent extends AbstractInitializable {
super(config);
this.transformManager = config.transformManager ?? getRoot().items.transformerManager;
this.cache = getRoot().items.cache();
this.db = getRoot().items.db();
this.componentRepo = new DrizzleComponentRepository(this.db, {logger: this.logger});
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(',')}`);
@@ -88,6 +86,8 @@ export default abstract class AbstractComponent extends AbstractInitializable {
name = this.name as string;
}
this.db = await getRoot().items.db();
this.componentRepo = new DrizzleComponentRepository(this.db, {logger: this.logger});
this.dbComponent = await this.componentRepo.findOrInsert({
mode: this.componentType,
type: this.type,
+8 -2
View File
@@ -16,11 +16,17 @@ import { SimpleError } from '../errors/MSErrors.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 => {
export const getDbPath = (name: string = 'msDb', workingDirectory?: string): string => {
if (isMemoryDb(name)) {
return MEMORY_DB_NAME;
}
return path.resolve(workingDirectory ?? configDir, `${name}.db`);
return path.resolve(workingDirectory ?? configDir, `${name}`);
}
export const getDbBackupPath = (dbPath: string, suffix?: string): string => {
const pathInfo = path.parse(dbPath);
const backupPath = `${path.join(pathInfo.dir, pathInfo.name)}.bak${suffix !== undefined ? `.${suffix}` : ''}`;
return backupPath;
}
export const backupDb = async (dbName: string, opts: { logger?: Logger, workingDirectory?: string } = {}): Promise<void> => {
@@ -1,46 +1,70 @@
import { drizzle } from 'drizzle-orm/node-sqlite';
import { drizzle as drizzlePglite } from 'drizzle-orm/pglite';
import { migrate } from 'drizzle-orm/node-sqlite/migrator';
import { migrate as migratePglite } from 'drizzle-orm/pglite/migrator';
import { BaseSQLiteDatabase } from "drizzle-orm/sqlite-core";
import { PGlite, PGliteOptions } from '@electric-sql/pglite';
import { sql as dsl, LogWriter, Logger as DrizzleLogger } from 'drizzle-orm';
import * as fs from 'fs/promises';
import * as fsSync from 'fs';
import * as path from 'path';
import { backupDb, getDbPath, MEMORY_DB_NAME } from '../Database.js';
import { fileExists } from '../../../utils/FSUtils.js';
import { backupDb, getDbBackupPath, getDbPath, MEMORY_DB_NAME } from '../Database.js';
import { fileExists, fileOrDirectoryIsWriteable } from '../../../utils/FSUtils.js';
import { childLogger, Logger, LogLevel } from '@foxxmd/logging';
import { loggerNoop } from '../../MaybeLogger.js';
import { projectDir } from '../../index.js';
import { relations } from './schema/schema.js';
import { addToContext, executeQuery } from './logContext.js';
export async function shouldBackupDb(dbPath: string, opts: {logger?: Logger, migrationsFolder?: string} = {}): Promise<[boolean, string[]]> {
export async function shouldBackupDb(dbVal: string | DbConcrete, opts: {logger?: Logger, migrationsFolder?: string} = {}): Promise<[boolean, string[]]> {
const {
logger: parentLogger = loggerNoop,
migrationsFolder = path.resolve(projectDir, 'src/backend/common/database/drizzle/migrations')
} = opts;
const logger = childLogger(parentLogger, 'Migrations');
let db: DbConcrete;
logger.info(`Checking database at ${dbPath}`);
if (dbPath !== MEMORY_DB_NAME && !fileExists(dbPath)) {
logger.info(`No database exists!`);
return [false, []];
if(typeof dbVal === 'string') {
logger.info(`Checking for database at ${dbVal}`);
if (dbVal !== MEMORY_DB_NAME && !fileExists(dbVal)) {
logger.info(`No database exists, no backup needed.`);
return [false, []];
}
db = await getDb(dbVal);
} else {
db = dbVal;
}
const db = drizzle(dbPath);
// const db = drizzlePglite(dbPath);
try {
// Ensure the migrations table exists
// https://github.com/drizzle-team/drizzle-orm/issues/1953
const res = db.all(dsl`
SELECT count(*) FROM sqlite_master WHERE type='table' AND name='__drizzle_migrations';
const res = await db.execute(dsl`
SELECT EXISTS (
SELECT FROM
pg_tables
WHERE
schemaname = 'drizzle' AND
tablename = '__drizzle_migrations'
);
`);
if (res[0]['count(*)'] === 0) {
// const res3 = await db.execute(dsl`
// SELECT * FROM
// pg_tables;
// `);
if (res.rows[0].exists === false) {
logger.info(`Database exists but there is no __drizzle_migrations table??`);
return [true, []];
}
const dbMigrations = await db.all(dsl`SELECT id, hash, created_at, name, applied_at FROM "__drizzle_migrations" ORDER BY created_at DESC`);
const appliedMigrations = new Set(dbMigrations.map((m: any) => m.name));
const dbMigrations = await db.execute(dsl`SELECT id, hash, created_at, name, applied_at FROM drizzle.__drizzle_migrations ORDER BY created_at DESC`);
// @ts-ignore
const appliedMigrations = new Set(dbMigrations.rows.map((m: any) => m.name));
const allFiles = await fs.readdir(migrationsFolder);
const migrationFiles = allFiles
@@ -61,25 +85,40 @@ export async function shouldBackupDb(dbPath: string, opts: {logger?: Logger, mig
} catch (error) {
logger.error(new Error('Failed to get pending migrations', { cause: error }));
return [true, []];
} finally {
if(db.$client.isOpen) {
db.$client.close();
}
}
}
export const getDb = (dbName: string = 'ms', opts: { logger?: Logger, workingDirectory?: string } = {}) => {
export const getDb = async (dbVal: string | PGlite, opts: { logger?: Logger, backupPath?: string, loadDataDir?: Promise<Blob> } = {}) => {
const {
workingDirectory,
logger = loggerNoop,
backupPath,
loadDataDir
} = opts;
const dbPath = getDbPath(dbName, workingDirectory);
return drizzle(dbPath, {relations: relations, logger: createDrizzleLogger(logger)});
let client: PGlite;
if(typeof dbVal === 'string') {
const opts: PGliteOptions = {};
if(dbVal !== MEMORY_DB_NAME) {
opts.dataDir = dbVal;
if(backupPath !== undefined) {
opts.loadDataDir = new Blob([fsSync.readFileSync(backupPath)]);
}
}
// only load one
// but this could be for a memory db so don't put it in above if
if(loadDataDir !== undefined && backupDb === undefined) {
opts.loadDataDir = await loadDataDir
}
client = await PGlite.create(opts);
} else {
client = dbVal;
}
return drizzlePglite({relations: relations, logger: createDrizzleLogger(logger), client});
}
export type DbConcrete = ReturnType<typeof getDb>;
export type DbConcrete = Awaited<ReturnType<typeof getDb>>;
export const migrateDb = async (db: ReturnType<typeof drizzle>, opts: {logger?: Logger, migrationsFolder?: string} = {}) => {
export const migrateDb = async (db: DbConcrete, opts: {logger?: Logger, migrationsFolder?: string} = {}) => {
const {
migrationsFolder,
logger: parentLogger = loggerNoop
@@ -88,7 +127,7 @@ export const migrateDb = async (db: ReturnType<typeof drizzle>, opts: {logger?:
try {
logger.info('Starting migrations...');
await executeQuery('migrations', async () => migrate(db, { migrationsFolder: migrationsFolder ?? path.resolve(projectDir, 'src/backend/common/database/drizzle/migrations') }), logger, process.env.LOG_MIGRATION === 'true' ? true : 'error');
await executeQuery('migrations', async () => migratePglite(db, { migrationsFolder: migrationsFolder ?? path.resolve(projectDir, 'src/backend/common/database/drizzle/migrations') }), logger, process.env.LOG_MIGRATION === 'true' ? true : 'error');
logger.info('Migrations complete');
} catch (e) {
throw new Error('Failed to migrate database', { cause: e });
@@ -111,15 +150,69 @@ export const migrateDbSync = (db: ReturnType<typeof drizzle>, opts: {logger?: Lo
}
}
export const performDbMigrationWithBackup = async (dbName: string = 'ms', opts: { logger?: Logger, workingDirectory?: string, migrationsFolder?: string } = {}) => {
const dbPath = getDbPath(dbName, opts.workingDirectory);
export const getMigratedDb = async (dbPath: string, opts: { logger?: Logger, workingDirectory?: string, migrationsFolder?: string, backupPath?: string, loadDataDir?: Promise<Blob> } = {}): Promise<[DbConcrete, boolean]> => {
const {
logger = loggerNoop
} = opts;
let db: DbConcrete,
isNew = false,
hasPendingMigrations: boolean = true;
if (dbPath !== MEMORY_DB_NAME) {
try {
fileOrDirectoryIsWriteable(dbPath);
} catch (e) {
throw new Error('Database directory is not accessible', { cause: e });
}
const [shouldBackup, pendingMigrations] = await shouldBackupDb(dbPath, opts);
if(shouldBackup) {
await backupDb(dbName, opts);
const backupPath = getDbBackupPath(dbPath);
if (fileExists(dbPath)) {
db = await getDb(dbPath, opts);
const [shouldBackup, pendingMigrations] = await shouldBackupDb(db, opts);
if (shouldBackup) {
hasPendingMigrations = true;
await backupPgDb(db, dbPath, { logger: opts.logger });
}
} else if(fileExists(backupPath)) {
logger.info(`Detected no database, using backup to recreate db. Backup file: ${backupPath}`);
db = await getDb(dbPath, {...opts, backupPath});
const usedBackedPath = getDbBackupPath(dbPath, 'used');
logger.info(`Backup loaded! Renaming backup to indicate it has already been used, new path: ${usedBackedPath}`);
await fs.rename(backupPath, usedBackedPath);
} else {
logger.info('Detected no database, creating a new one...');
db = await getDb(dbPath, opts);
isNew = true;
}
} else {
logger.info('Detected in-memory database');
db = await getDb(dbPath, opts);
isNew = true;
}
if(hasPendingMigrations && dbPath !== MEMORY_DB_NAME) {
logger.info('TIP: Migrations may take some time, depending on the size of your database');
}
const db = getDb(dbName, opts);
await migrateDb(db, opts);
return [db, isNew];
}
export const backupPgDb = async (db: DbConcrete, dbPath: string, opts: { logger?: Logger } = {}): Promise<void> => {
const {
logger: parentLogger = loggerNoop,
} = opts;
const logger = childLogger(parentLogger, 'Backup');
const pathInfo = path.parse(dbPath);
// being extra sure there isn't a trailing slash
const backupPath = `${path.join(pathInfo.dir, pathInfo.name)}-${Date.now()}.bak`;
logger.info(`Backing up database before migrating => ${backupPath}`);
fs.writeFile(backupPath, Buffer.from(await (await db.$client.dumpDataDir()).arrayBuffer()));
//await fs.copyFile(dbPath, backupPath)
logger.info('Backed up!');
}
export const createDrizzleLogger = (parentLogger: Logger, opts: {level?: LogLevel} = {}): DrizzleLogger => {
@@ -1,83 +0,0 @@
CREATE TABLE `components` (
`id` integer PRIMARY KEY,
`uid` text(200) NOT NULL,
`mode` text NOT NULL,
`type` text(50) NOT NULL,
`name` text NOT NULL,
`countLive` integer DEFAULT 0 NOT NULL,
`countNonLive` integer DEFAULT 0 NOT NULL,
`createdAt` number
);
--> statement-breakpoint
CREATE TABLE `jobs` (
`id` integer PRIMARY KEY,
`componentFromId` integer NOT NULL,
`componentToId` integer NOT NULL,
`name` text(50) NOT NULL,
`status` text DEFAULT 'idle' NOT NULL,
`retries` integer DEFAULT 0 NOT NULL,
`error` text,
`transformOptions` text,
`initialParameters` text,
`cursor` text,
`total` integer,
`imported` integer DEFAULT 0 NOT NULL,
`scrobbled` integer DEFAULT 0 NOT NULL,
`createdAt` number NOT NULL,
`updatedAt` number NOT NULL,
`completedAt` number,
CONSTRAINT `fk_jobs_componentFromId_components_id_fk` FOREIGN KEY (`componentFromId`) REFERENCES `components`(`id`) ON UPDATE CASCADE ON DELETE CASCADE,
CONSTRAINT `fk_jobs_componentToId_components_id_fk` FOREIGN KEY (`componentToId`) REFERENCES `components`(`id`) ON UPDATE CASCADE ON DELETE CASCADE
);
--> statement-breakpoint
CREATE TABLE `play_inputs` (
`id` integer PRIMARY KEY,
`playId` integer NOT NULL,
`data` text,
`play` text NOT NULL,
`createdAt` number,
CONSTRAINT `fk_play_inputs_playId_plays_id_fk` FOREIGN KEY (`playId`) REFERENCES `plays`(`id`) ON UPDATE CASCADE ON DELETE CASCADE
);
--> statement-breakpoint
CREATE TABLE `plays` (
`id` integer PRIMARY KEY,
`uid` text(30) NOT NULL,
`componentId` integer,
`error` text,
`playedAt` number,
`seenAt` number,
`updatedAt` number NOT NULL,
`play` text NOT NULL,
`state` text NOT NULL,
`parentId` integer,
`jobId` integer,
`playHash` text,
`mbidIdentifier` text,
`compacted` text,
CONSTRAINT `fk_plays_componentId_components_id_fk` FOREIGN KEY (`componentId`) REFERENCES `components`(`id`) ON UPDATE CASCADE ON DELETE CASCADE,
CONSTRAINT `fk_plays_parentId_plays_id_fk` FOREIGN KEY (`parentId`) REFERENCES `plays`(`id`) ON UPDATE CASCADE ON DELETE SET NULL,
CONSTRAINT `fk_plays_jobId_jobs_id_fk` FOREIGN KEY (`jobId`) REFERENCES `jobs`(`id`) ON UPDATE CASCADE ON DELETE CASCADE
);
--> statement-breakpoint
CREATE TABLE `play_queue_states` (
`id` integer PRIMARY KEY,
`playId` integer NOT NULL,
`componentId` integer NOT NULL,
`queueName` text(50) NOT NULL,
`queueStatus` text DEFAULT 'queued' NOT NULL,
`retries` integer DEFAULT 0 NOT NULL,
`error` text,
`createdAt` number NOT NULL,
`updatedAt` number NOT NULL,
CONSTRAINT `fk_play_queue_states_playId_plays_id_fk` FOREIGN KEY (`playId`) REFERENCES `plays`(`id`) ON UPDATE CASCADE ON DELETE CASCADE,
CONSTRAINT `fk_play_queue_states_componentId_components_id_fk` FOREIGN KEY (`componentId`) REFERENCES `components`(`id`) ON UPDATE CASCADE ON DELETE CASCADE
);
--> statement-breakpoint
CREATE UNIQUE INDEX `uid_mode_type_idx` ON `components` (`uid`,`mode`,`type`);--> statement-breakpoint
CREATE UNIQUE INDEX `play_input_id_idx` ON `play_inputs` (`playId`);--> statement-breakpoint
CREATE INDEX `play_parent_id_idx` ON `plays` (`parentId`);--> statement-breakpoint
CREATE INDEX `play_component_id_idx` ON `plays` (`componentId`);--> statement-breakpoint
CREATE UNIQUE INDEX `play_uid_idx` ON `plays` (`uid`);--> statement-breakpoint
CREATE INDEX `play_playedAt_idx` ON `plays` (`playedAt`);--> statement-breakpoint
CREATE INDEX `play_seenAt_idx` ON `plays` (`seenAt`);--> statement-breakpoint
CREATE INDEX `play_queue_state_id_idx` ON `play_queue_states` (`playId`);
@@ -0,0 +1,83 @@
CREATE TABLE "components" (
"id" serial PRIMARY KEY,
"uid" varchar(200) NOT NULL,
"mode" varchar(15) NOT NULL,
"type" varchar(50) NOT NULL,
"name" varchar NOT NULL,
"countLive" integer DEFAULT 0 NOT NULL,
"countNonLive" integer DEFAULT 0 NOT NULL,
"createdAt" timestamp
);
--> statement-breakpoint
CREATE TABLE "jobs" (
"id" serial PRIMARY KEY,
"componentFromId" integer NOT NULL,
"componentToId" integer NOT NULL,
"name" varchar(200) NOT NULL,
"status" varchar(20) DEFAULT 'idle' NOT NULL,
"retries" integer DEFAULT 0 NOT NULL,
"error" json,
"transformOptions" json,
"initialParameters" json,
"cursor" json,
"total" integer,
"imported" integer DEFAULT 0 NOT NULL,
"scrobbled" integer DEFAULT 0 NOT NULL,
"createdAt" timestamp NOT NULL,
"updatedAt" timestamp NOT NULL,
"completedAt" timestamp
);
--> statement-breakpoint
CREATE TABLE "play_inputs" (
"id" serial PRIMARY KEY,
"playId" integer NOT NULL,
"data" json,
"play" jsonb NOT NULL,
"createdAt" timestamp
);
--> statement-breakpoint
CREATE TABLE "plays" (
"id" serial PRIMARY KEY,
"uid" varchar(30) NOT NULL UNIQUE,
"componentId" integer,
"error" json,
"playedAt" timestamp,
"seenAt" timestamp,
"updatedAt" timestamp NOT NULL,
"play" jsonb NOT NULL,
"state" varchar(20) NOT NULL,
"parentId" integer,
"jobId" integer,
"playHash" varchar(100),
"mbidIdentifier" varchar(100),
"compacted" varchar(30)
);
--> statement-breakpoint
CREATE TABLE "play_queue_states" (
"id" serial PRIMARY KEY,
"playId" integer NOT NULL,
"componentId" integer NOT NULL,
"queueName" varchar(50) NOT NULL,
"queueStatus" varchar(20) DEFAULT 'queued' NOT NULL,
"retries" integer DEFAULT 0 NOT NULL,
"error" json,
"createdAt" timestamp NOT NULL,
"updatedAt" timestamp NOT NULL
);
--> statement-breakpoint
CREATE UNIQUE INDEX "uid_mode_type_idx" ON "components" ("uid","mode","type");--> statement-breakpoint
CREATE UNIQUE INDEX "play_input_id_idx" ON "play_inputs" ("playId");--> statement-breakpoint
CREATE INDEX "play_parent_id_idx" ON "plays" ("parentId");--> statement-breakpoint
CREATE INDEX "play_component_id_idx" ON "plays" ("componentId");--> statement-breakpoint
CREATE UNIQUE INDEX "play_uid_idx" ON "plays" ("uid");--> statement-breakpoint
CREATE INDEX "play_playedAt_idx" ON "plays" ("playedAt");--> statement-breakpoint
CREATE INDEX "play_seenAt_idx" ON "plays" ("seenAt");--> statement-breakpoint
CREATE INDEX "play_queue_state_id_idx" ON "play_queue_states" ("playId");--> statement-breakpoint
ALTER TABLE "jobs" ADD CONSTRAINT "jobs_componentFromId_components_id_fkey" FOREIGN KEY ("componentFromId") REFERENCES "components"("id") ON DELETE CASCADE ON UPDATE CASCADE;--> statement-breakpoint
ALTER TABLE "jobs" ADD CONSTRAINT "jobs_componentToId_components_id_fkey" FOREIGN KEY ("componentToId") REFERENCES "components"("id") ON DELETE CASCADE ON UPDATE CASCADE;--> statement-breakpoint
ALTER TABLE "play_inputs" ADD CONSTRAINT "play_inputs_playId_plays_id_fkey" FOREIGN KEY ("playId") REFERENCES "plays"("id") ON DELETE CASCADE ON UPDATE CASCADE;--> statement-breakpoint
ALTER TABLE "plays" ADD CONSTRAINT "plays_componentId_components_id_fkey" FOREIGN KEY ("componentId") REFERENCES "components"("id") ON DELETE CASCADE ON UPDATE CASCADE;--> statement-breakpoint
ALTER TABLE "plays" ADD CONSTRAINT "plays_parentId_plays_id_fkey" FOREIGN KEY ("parentId") REFERENCES "plays"("id") ON DELETE SET NULL ON UPDATE CASCADE;--> statement-breakpoint
ALTER TABLE "plays" ADD CONSTRAINT "plays_jobId_jobs_id_fkey" FOREIGN KEY ("jobId") REFERENCES "jobs"("id") ON DELETE CASCADE ON UPDATE CASCADE;--> statement-breakpoint
ALTER TABLE "play_queue_states" ADD CONSTRAINT "play_queue_states_playId_plays_id_fkey" FOREIGN KEY ("playId") REFERENCES "plays"("id") ON DELETE CASCADE ON UPDATE CASCADE;--> statement-breakpoint
ALTER TABLE "play_queue_states" ADD CONSTRAINT "play_queue_states_componentId_components_id_fkey" FOREIGN KEY ("componentId") REFERENCES "components"("id") ON DELETE CASCADE ON UPDATE CASCADE;
@@ -1,5 +1,5 @@
import { childLogger, Logger } from "@foxxmd/logging";
import { getDb } from "../drizzleUtils.js";
import { DbConcrete } from "../drizzleUtils.js";
import { CompareOpKey } from "../drizzleTypes.js";
import { Dayjs } from "dayjs";
import { RelationsFieldFilter, eq, inArray } from "drizzle-orm";
@@ -41,10 +41,10 @@ export abstract class DrizzleBaseRepository<T extends TableName> {
displayName: string;
tableName: TableName;
table: ReturnType<typeof getConfigByTableName<T>>
db: ReturnType<typeof getDb>;
db: DbConcrete;
componentId?: number
constructor(db: ReturnType<typeof getDb>, tableName: TableName, displayName: string, opts: DrizzleRepositoryOpts = {}) {
constructor(db: DbConcrete, tableName: TableName, displayName: string, opts: DrizzleRepositoryOpts = {}) {
this.db = db;
this.displayName = displayName;
this.tableName = tableName;
@@ -1,13 +1,13 @@
import { Logger } from "drizzle-orm";
import { DrizzleBaseRepository, DrizzleRepositoryOpts } from "./BaseRepository.js";
import { getDb } from "../drizzleUtils.js";
import { DbConcrete } from "../drizzleUtils.js";
import { ComponentNew, ComponentSelect, FindWhere } from "../drizzleTypes.js";
import { components } from "../schema/schema.js";
import { generateComponentEntity } from "../entityUtils.js";
export class DrizzleComponentRepository extends DrizzleBaseRepository<'components'> {
constructor(db: ReturnType<typeof getDb>, opts: DrizzleRepositoryOpts = {}) {
constructor(db: DbConcrete, opts: DrizzleRepositoryOpts = {}) {
super(db, 'components', 'Component', opts);
}
@@ -1,5 +1,5 @@
import { childLogger, Logger, LoggerAppExtras } from "@foxxmd/logging";
import { DbConcrete, getDb, runTransaction } from "../drizzleUtils.js";
import { DbConcrete } from "../drizzleUtils.js";
import { loggerNoop } from "../../../MaybeLogger.js";
import { ErrorLike, PlayObject, TA_CLOSE, TA_DEFAULT_ACCURACY, TA_EXACT, TemporalAccuracy } from "../../../../../core/Atomic.js";
import { generateInputEntity, generatePlayEntity, PlayEntityOpts, hydratePlaySelect, PlayHydateOptions } from "../entityUtils.js";
@@ -69,7 +69,7 @@ export class DrizzlePlayRepository extends DrizzleBaseRepository<'plays'> {
protected getQueueNextPrepared?: ReturnType<typeof this.prepareGetQueueNext>
protected getQueuedScrobbleRangePrepared?: ReturnType<typeof this.prepareGetQueuedScrobbleRange>
constructor(db: ReturnType<typeof getDb>, opts: DrizzleRepositoryOpts = {}) {
constructor(db: DbConcrete, opts: DrizzleRepositoryOpts = {}) {
super(db, 'plays', 'Plays', opts);
}
@@ -94,7 +94,7 @@ export class DrizzlePlayRepository extends DrizzleBaseRepository<'plays'> {
} = opts;
let playRows: PlaySelect[];
await runTransaction(this.db, async () => {
await this.db.transaction(async (tx) => {
const entitiesData = entitiesOpts.map((data) => {
const {
@@ -105,7 +105,7 @@ export class DrizzlePlayRepository extends DrizzleBaseRepository<'plays'> {
return generatePlayEntity(play, { componentId: this.componentId, ...rest });
});
playRows = await this.db.insert(plays).values(entitiesData).returning();
playRows = await tx.insert(plays).values(entitiesData).returning();
const inputDatas = playRows.map((x, index) => {
const {
@@ -120,7 +120,7 @@ export class DrizzlePlayRepository extends DrizzleBaseRepository<'plays'> {
return generateInputEntity({ play: inputPlay, playId: x.id, ...restInput });
});
const inputRow = await this.db.insert(playInputs).values(inputDatas);
const inputRow = await tx.insert(playInputs).values(inputDatas);
});
@@ -1,12 +1,12 @@
import { eq, and, lte, inArray } from "drizzle-orm";
import { DrizzleBaseRepository, DrizzleRepositoryOpts } from "./BaseRepository.js";
import { getDb } from "../drizzleUtils.js";
import { DbConcrete } from "../drizzleUtils.js";
import { QueueStateSelect } from "../drizzleTypes.js";
import { queueStates } from "../schema/schema.js";
import { CLIENT_DEAD_QUEUE } from "../../../../../core/Atomic.js";
export class DrizzleQueueRepository extends DrizzleBaseRepository<'queueStates'> {
constructor(db: ReturnType<typeof getDb>, opts: DrizzleRepositoryOpts = {}) {
constructor(db: DbConcrete, opts: DrizzleRepositoryOpts = {}) {
super(db, 'queueStates', 'Queue', opts);
}
@@ -1,4 +1,4 @@
import { integer, sqliteTable, text, index, uniqueIndex, customType, AnySQLiteColumn } from "drizzle-orm/sqlite-core";
import { integer, serial as primaryInt, pgTable as table, text, varchar, json, index, uniqueIndex, customType, AnyPgColumn, timestamp } from "drizzle-orm/pg-core";
import { defineRelations } from 'drizzle-orm';
import dayjs, { Dayjs } from "dayjs";
import { nanoid } from "nanoid";
@@ -10,16 +10,16 @@ import { JobRangeCount, JobRangeTime } from "../../../infrastructure/Job.js";
const DayjsTimestamp = customType<
{
data: Dayjs;
driverData: number;
driverData: string;
}
>({
dataType() {
return 'number'
return 'timestamp'
},
toDriver(value: Dayjs): number {
return value.valueOf();
toDriver(value: Dayjs): string {
return value.toISOString();
},
fromDriver(value: number): Dayjs {
fromDriver(value: string): Dayjs {
return dayjs(value);
},
});
@@ -31,7 +31,7 @@ const PlayJson = customType<
}
>({
dataType() {
return 'text'
return 'jsonb'
},
toDriver(value: PlayObject): string {
const {
@@ -46,28 +46,28 @@ const PlayJson = customType<
} = value;
return JSON.stringify(rest);
},
fromDriver(value: string): PlayObject {
return asPlayCheap(JSON.parse(value));
fromDriver(value: any): PlayObject {
return asPlayCheap(value);
},
});
export const plays = sqliteTable("plays", {
id: integer().primaryKey(),
uid: text({ length: 30 }).notNull().unique().$defaultFn(() => nanoid(20)),
export const plays = table("plays", {
id: primaryInt().primaryKey(),
uid: varchar({ length: 30 }).notNull().unique().$defaultFn(() => nanoid(20)),
componentId: integer().references(() => components.id, {onDelete: 'cascade', onUpdate: 'cascade'}),
error: text({ mode: 'json' }).$type<ErrorLike>(),
error: json().$type<ErrorLike>(),
playedAt: DayjsTimestamp('playedAt'),
seenAt: DayjsTimestamp('seenAt'),
updatedAt: DayjsTimestamp('updatedAt').notNull().$defaultFn(() => dayjs()),
play: PlayJson('play').notNull(), // text({ mode: 'json' }).notNull().$type<PlayObject>(),
state: text({enum: ['queued','discovered','discarded','scrobbled','failed','duped']}).notNull(),
state: varchar({enum: ['queued','discovered','discarded','scrobbled','failed','duped'], length: 20}).notNull(),
// https://orm.drizzle.team/docs/indexes-constraints#foreign-key
parentId: integer().references((): AnySQLiteColumn => plays.id, {onDelete: 'set null', onUpdate: 'cascade'}),
parentId: integer().references((): AnyPgColumn => plays.id, {onDelete: 'set null', onUpdate: 'cascade'}),
jobId: integer().references(() => jobs.id, {onDelete: 'cascade', onUpdate: 'cascade'}),
playHash: text(),
mbidIdentifier: text(),
compacted: text()
playHash: varchar({length: 100}),
mbidIdentifier: varchar({length: 100}),
compacted: varchar({length: 30})
}, (table) => [
index("play_parent_id_idx").on(table.parentId),
index("play_component_id_idx").on(table.componentId),
@@ -76,10 +76,10 @@ export const plays = sqliteTable("plays", {
index("play_seenAt_idx").on(table.seenAt)
]);
export const playInputs = sqliteTable("play_inputs", {
id: integer({ mode: 'number' }).primaryKey(),
export const playInputs = table("play_inputs", {
id: primaryInt().primaryKey(),
playId: integer().notNull().references(() => plays.id, {onDelete: 'cascade', onUpdate: 'cascade'}),
data: text({ mode: 'json' }).$type<object>(),
data: json().$type<object>(),
play: PlayJson('play').notNull(),//text({ mode: 'json' }).notNull().$type<PlayObject>(),
createdAt: DayjsTimestamp('createdAt').$defaultFn(() => dayjs())
}, (table) => [
@@ -107,14 +107,14 @@ export const playInputs = sqliteTable("play_inputs", {
// }
// }));
export const queueStates = sqliteTable("play_queue_states", {
id: integer({ mode: 'number' }).primaryKey(),
export const queueStates = table("play_queue_states", {
id: primaryInt().primaryKey(),
playId: integer().notNull().references(() => plays.id, {onDelete: 'cascade', onUpdate: 'cascade'}),
componentId: integer().notNull().references(() => components.id, {onDelete: 'cascade', onUpdate: 'cascade'}),
queueName: text({length: 50}).notNull(),
queueStatus: text({enum: ['queued','completed','failed']}).notNull().default('queued'),
queueName: varchar({length: 50}).notNull(),
queueStatus: varchar({enum: ['queued','completed','failed'], length: 20}).notNull().default('queued'),
retries: integer().notNull().default(0),
error: text({ mode: 'json' }).$type<ErrorLike>(),
error: json().$type<ErrorLike>(),
createdAt: DayjsTimestamp('createdAt').notNull().$defaultFn(() => dayjs()),
updatedAt: DayjsTimestamp('updatedAt').notNull().$defaultFn(() => dayjs())
}, (table) => [
@@ -133,16 +133,16 @@ export const queueStates = sqliteTable("play_queue_states", {
// }
// }));
export const components = sqliteTable("components", {
id: integer({ mode: 'number' }).primaryKey(),
export const components = table("components", {
id: primaryInt().primaryKey(),
// user-provided id
uid: text({ length: 200 }).notNull(),
mode: text({enum: ['source','client']}).notNull(),
uid: varchar({ length: 200 }).notNull(),
mode: varchar({enum: ['source','client'], length: 15}).notNull(),
// spotify, lastfm, etc...
type: text({length: 50}).notNull(),
type: varchar({length: 50}).notNull(),
// vanity display name
// used as uid if no user-provided id
name: text().notNull(),
name: varchar().notNull(),
// number of discovered/scrobbled plays found in real time
countLive: integer().notNull().default(0),
// number of discovered/scrobbled plays from backlog/jobs
@@ -153,17 +153,17 @@ export const components = sqliteTable("components", {
uniqueIndex('uid_mode_type_idx').on(table.uid,table.mode,table.type)
]);
export const jobs = sqliteTable("jobs", {
id: integer({ mode: 'number' }).primaryKey(),
export const jobs = table("jobs", {
id: primaryInt().primaryKey(),
componentFromId: integer().notNull().references(() => components.id, {onDelete: 'cascade', onUpdate: 'cascade'}),
componentToId: integer().notNull().references(() => components.id, {onDelete: 'cascade', onUpdate: 'cascade'}),
name: text({length: 50}).notNull(),
status: text({enum: ['idle','completed','failed','processing']}).notNull().default('idle'),
name: varchar({length: 200}).notNull(),
status: varchar({enum: ['idle','completed','failed','processing'], length: 20}).notNull().default('idle'),
retries: integer().notNull().default(0),
error: text({ mode: 'json' }).$type<ErrorLike>(),
transformOptions: text({ mode: 'json' }).$type<PlayTransformPartsConfig<SearchAndReplaceTerm[] | ExternalMetadataTerm>>(),
initialParameters: text({ mode: 'json' }).$type<JobRangeCount | JobRangeTime>(),
cursor: text({ mode: 'json' }),
error: json().$type<ErrorLike>(),
transformOptions: json().$type<PlayTransformPartsConfig<SearchAndReplaceTerm[] | ExternalMetadataTerm>>(),
initialParameters: json().$type<JobRangeCount | JobRangeTime>(),
cursor: json(),
total: integer(),
imported: integer().notNull().default(0),
scrobbled: integer().notNull().default(0),
+3 -1
View File
@@ -438,4 +438,6 @@ export const REFRESH_STALE_DEFAULT = 60;
*
* @example [60, 3600, "1 hour", "4 days"]
*/
export type DurationValue = number | string;
export type DurationValue = number | string;
export type DbExternalMode = 'none' | 'live' | 'standalone';
+194 -109
View File
@@ -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 { getDb, performDbMigrationWithBackup } 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,159 @@ 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}`);
logger.info(`Using database at ${getDbPath('ms')}`);
await performDbMigrationWithBackup('ms', {logger});
const dbPath = getDbPath('msDb');
logger.info(`Using database at ${getDbPath('msDb')}`);
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: getDb('ms', {logger})
});
const internalConfigOptional = {
localUrl: root.get('localUrl'),
configDir: root.get('configDir'),
version: root.get('version')
};
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'));
await root.items.cache().init(true);
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');
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');
}
if(appConfigFail !== undefined) {
logger.warn('App config file exists but could not be parsed!');
logger.warn(appConfigFail);
}
if (dbMode === 'standalone') {
logger.info('MS App startup stopped early due to Standalone DB Mode.');
} else {
const notifiers = new Notifiers(root.get('notifierEmitter'), root.get('clientEmitter'), root.get('sourceEmitter'), root.get('logger')); //root.get('notifiers');
await notifiers.buildWebhooks(webhooks);
const root = getRoot({
...config,
cache: parseUserConfig(cache, logger),
logger,
loggingConfig: logging,
loggerStream: appLoggerStream,
db
});
await root.items.transformerManager.registerFromEnv();
await root.items.transformerManager.registeryDefaults();
await root.items.transformerManager.initTransformers();
const internalConfigOptional = {
localUrl: root.get('localUrl'),
configDir: root.get('configDir'),
version: root.get('version')
};
/*
* setup clients
* */
await scrobbleClients.buildClientsFromConfig(notifiers);
/*
* setup sources
* */
await scrobbleSources.buildSourcesFromConfig([]);
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'));
// check ambiguous client/source types like this for now
const lastfmSources = scrobbleSources.getByType('lastfm');
const lastfmScrobbles = scrobbleClients.getByType('lastfm');
await root.items.cache().init(true);
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');
}
if (appConfigFail !== undefined) {
logger.warn('App config file exists but could not be parsed!');
logger.warn(appConfigFail);
}
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...`);
}
}
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();
}
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();
}
}
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...`);
}
}
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();
}
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});
if(logger !== undefined) {
@@ -223,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);
}
}());
+4 -5
View File
@@ -28,7 +28,7 @@ export interface RootOptions {
cache?: CacheConfigOptions | MSCache | (() => MSCache)
mbMap?: MusicBrainzSingletonMap | (() => MusicBrainzSingletonMap)
transformers?: TransformerCommonConfig[]
db?: DbConcrete | (() => DbConcrete)
db?: DbConcrete | (() => Promise<DbConcrete>)
}
const discovered = new prom.Counter({
@@ -92,12 +92,11 @@ const createRoot = (options: RootOptions = {logger: loggerDebug}) => {
maybeSingletonMb = new Map();
}
let dbFunc: () => DbConcrete;
let maybeSingletonDb: DbConcrete;
let dbFunc: () => Promise<DbConcrete>;
if(typeof db === 'function') {
dbFunc = db;
} else {
maybeSingletonDb = db;
dbFunc = async () => db;
}
const cEmitter = new WildcardEmitter();
@@ -158,7 +157,7 @@ const createRoot = (options: RootOptions = {logger: loggerDebug}) => {
cache: () => maybeSingletonCache !== undefined ? () => maybeSingletonCache : cacheFunc,
mbMap: () => maybeSingletonMb !== undefined ? () => maybeSingletonMb : mbFunc,
coverArtApi,
db: () => maybeSingletonDb !== undefined ? () => maybeSingletonDb : dbFunc
db: () => dbFunc
}).add((items) => {
const localUrl = generateBaseURL(baseUrl, items.port)
return {
@@ -152,8 +152,8 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i
declare protected componentType: 'client';
protected playRepo: DrizzlePlayRepository;
protected queueRepo: DrizzleQueueRepository;
protected playRepo!: DrizzlePlayRepository;
protected queueRepo!: DrizzleQueueRepository;
constructor(type: any, name: any, config: CommonClientConfig, notifier: Notifiers, emitter: EventEmitter, logger: Logger) {
super(config);
@@ -166,8 +166,6 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i
this.deadLogger = childLogger(this.logger, CLIENT_DEAD_QUEUE);
this.notifier = notifier;
this.emitter = emitter;
this.playRepo = new DrizzlePlayRepository(this.db, {logger: this.logger,});
this.queueRepo = new DrizzleQueueRepository(this.db, {logger: this.logger});
const {
options: {
@@ -327,6 +325,8 @@ export default abstract class AbstractScrobbleClient extends AbstractComponent i
}
protected async postDatabase(): Promise<void> {
this.playRepo = new DrizzlePlayRepository(this.db, {logger: this.logger});
this.queueRepo = new DrizzleQueueRepository(this.db, {logger: this.logger});
this.playRepo.componentId = this.dbComponent.id;
this.queueRepo.componentId = this.dbComponent.id;
this.tracksScrobbled = this.dbComponent.countLive + this.dbComponent.countNonLive;
+2 -2
View File
@@ -109,7 +109,7 @@ export default abstract class AbstractSource extends AbstractComponent implement
declare protected componentType: 'source';
protected playRepo: DrizzlePlayRepository;
protected playRepo!: DrizzlePlayRepository;
existingDiscoveredPlay: (playObjPre: PlayObject, existingScrobbles: PlayObject[], log?: boolean) => Promise<PlayMatchResult>
@@ -130,7 +130,6 @@ export default abstract class AbstractSource extends AbstractComponent implement
this.emitter = emitter;
this.discoveredCounter = getRoot().items.sourceMetics.discovered;
this.playRepo = new DrizzlePlayRepository(this.db, {logger: this.logger});
const existingScrobbleOpts: ExistingScrobbleOpts = {
logger: this.logger,
@@ -207,6 +206,7 @@ export default abstract class AbstractSource extends AbstractComponent implement
}
protected async postDatabase(): Promise<void> {
this.playRepo = new DrizzlePlayRepository(this.db, {logger: this.logger});
this.tracksDiscovered = this.dbComponent.countLive;
this.playRepo.componentId = this.dbComponent.id;
}
+116 -79
View File
@@ -1,6 +1,6 @@
import chai, { assert, expect } from 'chai';
import asPromised from 'chai-as-promised';
import { getDb, migrateDb, performDbMigrationWithBackup, shouldBackupDb } from '../../common/database/drizzle/drizzleUtils.js';
import { getDb, getMigratedDb, migrateDb, shouldBackupDb } from '../../common/database/drizzle/drizzleUtils.js';
import withLocalTmpDir from 'with-local-tmp-dir';
import { components, playInputs, plays, queueStates } from '../../common/database/drizzle/schema/schema.js';
import dayjs from 'dayjs';
@@ -11,14 +11,16 @@ import * as path from 'path';
import * as fs from 'fs/promises';
import { projectDir } from '../../common/index.js';
import { DatabaseSync } from 'node:sqlite';
import { fixtureCreateComponent, fixtureCreateInput, fixtureCreatePlay } from '../utils/databaseFixtures.js';
import { fixtureCreateComponent, fixtureCreateInput, fixtureCreatePlay, getPrepopulatedFSPGlite, getPrepopulatedMemoryPGlite } from '../utils/databaseFixtures.js';
import { DrizzlePlayRepository, RepositoryCreatePlayOpts } from '../../common/database/drizzle/repositories/PlayRepository.js';
import { generatePlayWithLifecycle, generateRandomObj } from '../../../core/tests/utils/fixtures.js';
import { generateArray } from '../../../core/DataUtils.js';
import { formatNumber, generateArray } from '../../../core/DataUtils.js';
import { objectsEqual } from '../../utils/DataUtils.js';
import { eq, sql } from 'drizzle-orm';
import { PlaySelect } from '../../common/database/drizzle/drizzleTypes.js';
import { loggerDebug } from '@foxxmd/logging';
import { transientDb } from '../utils/TransientTestUtils.js';
import { dataDir } from '@electric-sql/pglite-prepopulatedfs'
// would be great to push migrations directly from schema but doesn't seem supported in newest beta
// https://github.com/drizzle-team/drizzle-orm/discussions/4373
@@ -27,28 +29,32 @@ describe('Migrations', function () {
it('Detects non-existent db', async function () {
this.timeout(5000);
await withLocalTmpDir(async () => {
const [shouldBackup, pending] = await shouldBackupDb(getDbPath('notreal', process.cwd()));
expect(shouldBackup).is.false;
expect(pending).length(0);
const [db, isNew] = await getMigratedDb(getDbPath('notreal', process.cwd()), { loadDataDir: dataDir() });
expect(isNew).is.true;
db.$client.close();
}, {postfix: 'noDb'});
});
it('Detects abnormal db', async function () {
this.timeout(5000);
await withLocalTmpDir(async () => {
const otherdb = new DatabaseSync(path.resolve('./', 'other.db'));
const [shouldBackup, pending] = await shouldBackupDb(getDbPath('other', process.cwd()));
expect(shouldBackup).is.true;
expect(pending).length(0);
otherdb.close();
}, { unsafeCleanup: true, postfix: 'badDb' });
// database exists but there is no __drizzle_migrations table
const db = await getDb(':memory:', { loadDataDir: dataDir() });
const [shouldBackup, pending] = await shouldBackupDb(db);
expect(shouldBackup).is.true;
db.$client.close();
expect(pending).length(0);
});
it('Detects pending migrations', async function () {
this.timeout(5000);
const allFiles = await fs.readdir(path.resolve(projectDir, 'src/backend/common/database/drizzle/migrations'));
const migrationFiles = allFiles
.sort();
@@ -60,7 +66,7 @@ describe('Migrations', function () {
try {
await fs.cp(path.resolve(projectDir, `src/backend/common/database/drizzle/migrations/${migrationFiles[0]}`), path.resolve('./migrations/', migrationFiles[0]), { recursive: true });
const mf = path.resolve('./migrations');
const db = getDb('ms', { workingDirectory: process.cwd() });
const db = await getDb(':memory:', {loadDataDir: dataDir()});
await migrateDb(db, { migrationsFolder: mf });
const res = await x('drizzle-kit', [
'generate',
@@ -72,9 +78,9 @@ describe('Migrations', function () {
'--schema',
path.resolve(projectDir, 'src/backend/common/database/drizzle/schema'),
'--dialect',
'sqlite'
]);
const [shouldBackup, pending] = await shouldBackupDb(getDbPath('ms', process.cwd()), { migrationsFolder: mf });
'postgresql'
], {throwOnError: true});
const [shouldBackup, pending] = await shouldBackupDb(db, { migrationsFolder: mf });
expect(shouldBackup).is.true;
expect(pending).length(1);
expect(pending[0]).includes('newMigration');
@@ -87,6 +93,8 @@ describe('Migrations', function () {
it('Detects no pending migrations correctly', async function () {
this.timeout(5000);
const allFiles = await fs.readdir(path.resolve(projectDir, 'src/backend/common/database/drizzle/migrations'));
const migrationFiles = allFiles
.sort();
@@ -98,9 +106,9 @@ describe('Migrations', function () {
try {
await fs.cp(path.resolve(projectDir, `src/backend/common/database/drizzle/migrations/${migrationFiles[0]}`), path.resolve('./migrations/', migrationFiles[0]), { recursive: true });
const mf = path.resolve('./migrations');
const db = getDb('ms', { workingDirectory: process.cwd() });
const db = await getDb(':memory:', {loadDataDir: dataDir()});
await migrateDb(db, { migrationsFolder: mf });
const [shouldBackup, pending] = await shouldBackupDb(getDbPath('ms', process.cwd()), { migrationsFolder: mf });
const [shouldBackup, pending] = await shouldBackupDb(db, { migrationsFolder: mf });
expect(shouldBackup).is.false;
expect(pending).length(0);
db.$client.close();
@@ -112,19 +120,23 @@ describe('Migrations', function () {
it('Backs up database when migrations are pending', async function () {
// this can be slow due to all the io
this.timeout(5000);
const allFiles = await fs.readdir(path.resolve(projectDir, 'src/backend/common/database/drizzle/migrations'));
const migrationFiles = allFiles
.sort();
await withLocalTmpDir(async () => {
const dbPath = getDbPath('msDb', process.cwd());
// copy first migration
await fs.mkdir('migrations');
try {
await fs.cp(path.resolve(projectDir, `src/backend/common/database/drizzle/migrations/${migrationFiles[0]}`), path.resolve('./migrations/', migrationFiles[0]), { recursive: true });
const mf = path.resolve('./migrations');
const db = getDb('ms', { workingDirectory: process.cwd() });
await migrateDb(db, { migrationsFolder: mf });
const [db, _] = await getMigratedDb(dbPath, {migrationsFolder: mf, loadDataDir: dataDir()});
await db.$client.close();
const res = await x('drizzle-kit', [
'generate',
'--name',
@@ -135,17 +147,17 @@ describe('Migrations', function () {
'--schema',
path.resolve(projectDir, 'src/backend/common/database/drizzle/schema'),
'--dialect',
'sqlite'
]);
db.$client.close();
'postgresql'
], {throwOnError: true});
// add dummy data to migration so migrate() doesn't fail
const newMigrationFolder = (await fs.readdir(path.resolve('./migrations/'))).find(x => x.includes('newMigration'));
await fs.appendFile(path.resolve('./migrations/',newMigrationFolder, 'migration.sql'),`\nselect count(*) from plays;`);
await performDbMigrationWithBackup('ms', {workingDirectory: process.cwd(), migrationsFolder: mf});
await getMigratedDb(dbPath, {migrationsFolder: mf});
const contents = await fs.readdir(path.resolve('./'));
expect(contents.some(x => x.includes('ms.db.bak')));
const backupPattern = new RegExp(/msDb-\d+\.bak/)
expect(contents.some(x => backupPattern.test(x))).is.true;
} catch (e) {
throw e;
}
@@ -158,8 +170,7 @@ describe('Basic DB Operations', function () {
it('Should create a play', async function () {
const db = getDb(':memory:', { workingDirectory: process.cwd() });
await migrateDb(db);
const db = await transientDb();
const component = await db.insert(components).values(fixtureCreateComponent()).returning();
@@ -169,16 +180,15 @@ describe('Basic DB Operations', function () {
playedAt: dayjs(),
seenAt: dayjs(),
play: generatePlay()
});
}).returning();
expect(playRow.changes).eq(1);
expect(playRow.length).eq(1);
db.$client.close();
});
it('Should create a play with relations', async function () {
const db = getDb(':memory:', { workingDirectory: process.cwd() });
await migrateDb(db);
const db = await transientDb();
try {
@@ -226,8 +236,7 @@ describe('Basic DB Operations', function () {
it('deletes all dependent relations when a Play is deleted', async function () {
const db = getDb(':memory:', { workingDirectory: process.cwd() });
await migrateDb(db);
const db = await transientDb();
try {
@@ -301,8 +310,7 @@ describe('Repository Operations', function () {
it('creates Plays and inputs', async function () {
const db = getDb(':memory:');
await migrateDb(db);
const db = await transientDb();
const component = await db.insert(components).values(fixtureCreateComponent()).returning();
@@ -331,7 +339,7 @@ describe('Repository Operations', function () {
it('finds Plays by state', async function () {
const db = getDb(':memory:');
const db = await transientDb();
await migrateDb(db);
const component = await db.insert(components).values(fixtureCreateComponent()).returning();
@@ -361,7 +369,7 @@ describe('Repository Operations', function () {
it('finds Plays by date range', async function () {
const db = getDb(':memory:');
const db = await transientDb();
await migrateDb(db);
const component = await db.insert(components).values(fixtureCreateComponent()).returning();
@@ -410,7 +418,7 @@ describe('Repository Operations', function () {
it('finds Plays by component', async function () {
const db = getDb(':memory:');
const db = await transientDb();
await migrateDb(db);
const component1 = await db.insert(components).values(fixtureCreateComponent()).returning();
@@ -457,7 +465,7 @@ describe('Repository Operations', function () {
it('finds purgable Plays', async function () {
const db = getDb(':memory:');
const db = await transientDb();
await migrateDb(db);
const component1 = await db.insert(components).values(fixtureCreateComponent()).returning();
@@ -521,34 +529,34 @@ describe('Repository Operations', function () {
expect(p2Plays[1]).to.eq(childPlays[0].id);
});
it('Get json property from play', async function () {
// it('Get json property from play', async function () {
const db = getDb(':memory:', { workingDirectory: process.cwd() });
await migrateDb(db);
// const db = getDb(':memory:', { workingDirectory: process.cwd() });
// await migrateDb(db);
try {
// try {
const component = await db.insert(components).values(fixtureCreateComponent()).returning();
// const component = await db.insert(components).values(fixtureCreateComponent()).returning();
const playRows = await db.insert(plays).values([
fixtureCreatePlay({ componentId: component[0].id, play: generatePlay({}, {source: 'test1'}) }),
fixtureCreatePlay({ componentId: component[0].id, play: generatePlay({}, {source: 'test2'}) })
]).returning();
// const playRows = await db.insert(plays).values([
// fixtureCreatePlay({ componentId: component[0].id, play: generatePlay({}, {source: 'test1'}) }),
// fixtureCreatePlay({ componentId: component[0].id, play: generatePlay({}, {source: 'test2'}) })
// ]).returning();
let result: PlaySelect[];
// https://github.com/drizzle-team/drizzle-orm/discussions/938#discussioncomment-6542336
result = await db.select().from(plays).where(
sql`json_extract(${plays.play}, '$.meta.source') = 'test1'`
);
// let result: PlaySelect[];
// // https://github.com/drizzle-team/drizzle-orm/discussions/938#discussioncomment-6542336
// result = await db.select().from(plays).where(
// sql`json_extract(${plays.play}, '$.meta.source') = 'test1'`
// );
expect(result).length(1);
expect(result[0].play.meta.source).eq('test1');
// expect(result).length(1);
// expect(result[0].play.meta.source).eq('test1');
} catch (e) {
throw e;
}
db.$client.close();
});
// } catch (e) {
// throw e;
// }
// db.$client.close();
// });
});
@@ -565,10 +573,10 @@ describe('DB Size Stats', function () {
await withLocalTmpDir(async () => {
try {
let db = getDb('ms', { workingDirectory: process.cwd() });
let db = await getDb(await getPrepopulatedFSPGlite(getDbPath('msDb', process.cwd())));
await migrateDb(db);
const stats = await fs.stat(path.resolve('./ms.db'));
loggerDebug.debug(`Empty => ${stats.size / 1024}kb`);
const play100Component = Number.parseInt((await x('du', ['-ksb','.'])).stdout.split('\t')[0]);
loggerDebug.debug(`100 Plays => ${formatNumber((play100Component / 1024) / 1024, {toFixed: 2})}mb`);
} catch (e) {
throw e;
}
@@ -577,9 +585,11 @@ describe('DB Size Stats', function () {
it('get db plays size stats', async function () {
this.timeout(10000);
await withLocalTmpDir(async () => {
try {
let db = getDb('ms', { workingDirectory: process.cwd() });
let db = await getDb(await getPrepopulatedFSPGlite(getDbPath('msDb', process.cwd())));
await migrateDb(db);
const component = await db.insert(components).values(fixtureCreateComponent()).returning();
@@ -587,13 +597,21 @@ describe('DB Size Stats', function () {
const playData = generateArray<RepositoryCreatePlayOpts>(100, () => ({ ...fixtureCreatePlay({ componentId: component[0].id, play: generatePlay() }), state: 'queued', input: { data: undefined } }));
await playRepo.createPlays(playData);
const Play100Component = await fs.stat(path.resolve('./ms.db'));
loggerDebug.debug(`100 Plays => ${Play100Component.size / 1024}kb`);
const play100Component = Number.parseInt((await x('du', ['-ksb','.'])).stdout.split('\t')[0]);
loggerDebug.debug(`100 Plays => ${formatNumber((play100Component / 1024) / 1024, {toFixed: 2})}mb`);
const morePlayData = generateArray<RepositoryCreatePlayOpts>(900, () => ({ ...fixtureCreatePlay({ componentId: component[0].id, play: generatePlay() }), state: 'queued', input: { data: undefined }}));
await playRepo.createPlays(morePlayData);
const Play1000Component = await fs.stat(path.resolve('./ms.db'));
loggerDebug.debug(`1000 Plays => ${Play1000Component.size / 1024}kb`);
const play1000Component = Number.parseInt((await x('du', ['-ksb','.'])).stdout.split('\t')[0]);
loggerDebug.debug(`1000 Plays => ${formatNumber((play1000Component / 1024) / 1024, {toFixed: 2})}mb`);
for(let i = 0; i < 9; i++) {
const evenMorePlayData = generateArray<RepositoryCreatePlayOpts>(1000, () => ({ ...fixtureCreatePlay({ componentId: component[0].id, play: generatePlay() }), state: 'queued', input: { data: undefined }}));
await playRepo.createPlays(evenMorePlayData);
}
const play10000Component = Number.parseInt((await x('du', ['-ksb','.'])).stdout.split('\t')[0]);
loggerDebug.debug(`10000 Plays => ${formatNumber((play10000Component / 1024) / 1024, {toFixed: 2})}mb`);
} catch (e) {
throw e;
}
@@ -602,9 +620,11 @@ describe('DB Size Stats', function () {
it('get db plays size stats with input', async function () {
this.timeout(10000);
await withLocalTmpDir(async () => {
try {
let db = getDb('ms', { workingDirectory: process.cwd() });
let db = await getDb(await getPrepopulatedFSPGlite(getDbPath('msDb', process.cwd())));
await migrateDb(db);
const component = await db.insert(components).values(fixtureCreateComponent()).returning();
@@ -612,13 +632,20 @@ describe('DB Size Stats', function () {
const playData = generateArray<RepositoryCreatePlayOpts>(100, () => ({ ...fixtureCreatePlay({ componentId: component[0].id, play: generatePlay() }), state: 'queued', input: { data: generateRandomObj(undefined, { allowUndefined: false }) } }));
await playRepo.createPlays(playData);
const Play100Component = await fs.stat(path.resolve('./ms.db'));
loggerDebug.debug(`100 Plays => ${Play100Component.size / 1024}kb`);
const play100Component = Number.parseInt((await x('du', ['-ksb','.'])).stdout.split('\t')[0]);
loggerDebug.debug(`100 Plays => ${formatNumber((play100Component / 1024) / 1024, {toFixed: 2})}mb`);
const morePlayData = generateArray<RepositoryCreatePlayOpts>(900, () => ({ ...fixtureCreatePlay({ componentId: component[0].id, play: generatePlay() }), state: 'queued', input: { data: generateRandomObj(undefined, { allowUndefined: false }) } }));
await playRepo.createPlays(morePlayData);
const Play1000Component = await fs.stat(path.resolve('./ms.db'));
loggerDebug.debug(`1000 Plays => ${Play1000Component.size / 1024}kb`);
const play1000Component = Number.parseInt((await x('du', ['-ksb','.'])).stdout.split('\t')[0]);
loggerDebug.debug(`1000 Plays => ${formatNumber((play1000Component / 1024) / 1024, {toFixed: 2})}mb`);
for(let i = 0; i < 9; i++) {
const evenMorePlayData = generateArray<RepositoryCreatePlayOpts>(1000, () => ({ ...fixtureCreatePlay({ componentId: component[0].id, play: generatePlay() }), state: 'queued', input: { data: generateRandomObj(undefined, { allowUndefined: false }) } }));
await playRepo.createPlays(evenMorePlayData);
}
const play10000Component = Number.parseInt((await x('du', ['-ksb','.'])).stdout.split('\t')[0]);
loggerDebug.debug(`10000 Plays => ${formatNumber((play10000Component / 1024) / 1024, {toFixed: 2})}mb`);
} catch (e) {
throw e;
}
@@ -627,9 +654,11 @@ describe('DB Size Stats', function () {
it('get db plays size stats with input and lifecycle', async function () {
this.timeout(10000);
await withLocalTmpDir(async () => {
try {
let db = getDb('ms', { workingDirectory: process.cwd() });
let db = await getDb(await getPrepopulatedFSPGlite(getDbPath('msDb', process.cwd())));
await migrateDb(db);
const component = await db.insert(components).values(fixtureCreateComponent()).returning();
@@ -637,13 +666,21 @@ describe('DB Size Stats', function () {
const playData = generateArray<RepositoryCreatePlayOpts>(100, () => ({ ...fixtureCreatePlay({ componentId: component[0].id, play: generatePlayWithLifecycle({lifecycleSteps: {preCompare: 1}}) }), state: 'queued', input: { data: generateRandomObj(undefined, { allowUndefined: false }) } }));
await playRepo.createPlays(playData);
const Play100Component = await fs.stat(path.resolve('./ms.db'));
loggerDebug.debug(`100 Plays => ${Play100Component.size / 1024}kb`);
const play100Component = Number.parseInt((await x('du', ['-ksb','.'])).stdout.split('\t')[0]);
loggerDebug.debug(`100 Plays => ${formatNumber((play100Component / 1024) / 1024, {toFixed: 2})}mb`);
const morePlayData = generateArray<RepositoryCreatePlayOpts>(900, () => ({ ...fixtureCreatePlay({ componentId: component[0].id, play: generatePlayWithLifecycle({lifecycleSteps: {preCompare: 1}}) }), state: 'queued', input: { data: generateRandomObj(undefined, { allowUndefined: false }) } }));
await playRepo.createPlays(morePlayData);
const Play1000Component = await fs.stat(path.resolve('./ms.db'));
loggerDebug.debug(`1000 Plays => ${Play1000Component.size / 1024}kb`);
const play1000Component = Number.parseInt((await x('du', ['-ksb','.'])).stdout.split('\t')[0]);
loggerDebug.debug(`1000 Plays => ${formatNumber((play1000Component / 1024) / 1024, {toFixed: 2})}mb`);
for(let i = 0; i < 9; i++) {
const evenMorePlayData = generateArray<RepositoryCreatePlayOpts>(1000, () => ({ ...fixtureCreatePlay({ componentId: component[0].id, play: generatePlayWithLifecycle({lifecycleSteps: {preCompare: 1}}) }), state: 'queued', input: { data: generateRandomObj(undefined, { allowUndefined: false }) } }));
await playRepo.createPlays(evenMorePlayData);
}
const play10000Component = Number.parseInt((await x('du', ['-ksb','.'])).stdout.split('\t')[0]);
loggerDebug.debug(`10000 Plays => ${formatNumber((play10000Component / 1024) / 1024, {toFixed: 2})}mb`);
} catch (e) {
throw e;
}
+6 -2
View File
@@ -31,8 +31,6 @@ export class TestScrobbler extends AbstractScrobbleClient {
this.scrobbleDelay = 10;
this.scrobbleSleep = 20;
this.scrobbleWaitStopInterval = 20;
this.playRepoTest = this.playRepo;
this.queueRepoTest = this.queueRepo;
}
doScrobble(playObj: PlayObject) {
@@ -44,6 +42,12 @@ export class TestScrobbler extends AbstractScrobbleClient {
return super.doParseCache();
}
protected async postDatabase(): Promise<void> {
super.postDatabase();
this.playRepoTest = this.playRepo;
this.queueRepoTest = this.queueRepo;
}
playToClientPayload(playObject: PlayObject): object {
return playObject;
}
@@ -621,7 +621,10 @@ describe('Dead Scrobbles', function() {
}
await testScrobbler.processDeadLetterQueue();
await pEvent(testScrobbler.emitter, 'queueState');
await Promise.race([
sleep(15000),
pEvent(testScrobbler.emitter, 'queueState')
])
expect(testScrobbler.deadLetterQueued).eq(0);
});
+13
View File
@@ -2,5 +2,18 @@ import { loggerTest } from '@foxxmd/logging';
import { getRoot } from "../ioc.js";
import { transientCache, transientDb } from './utils/TransientTestUtils.js';
// let transientD: DbConcrete;
// const transientDbFactory = () => {
// return getDb(transientD.$client.clone())
// }
// export async function mochaGlobalSetup() {
// transientD = getDb(':memory:');
// await migrateDb(transientD);
// const root = getRoot({cache: transientCache, logger: loggerTest, db: transientDb});
// root.items.cache().init();
// }
const root = getRoot({cache: transientCache, logger: loggerTest, db: transientDb});
root.items.cache().init();
+10 -1
View File
@@ -426,8 +426,17 @@ describe('Player Cleanup', function () {
});
});
class DeezerTestSource extends DeezerInternalSource {
protected async doCheckConnection(): Promise<true | string | undefined> {
return;
}
doAuthentication = async () => {
return true;
}
}
const generateDeezerSource = async (options: DeezerInternalSourceOptions = {}) => {
const source = new DeezerInternalSource('test', {data: {arl: 'test'}, options}, {localUrl: new URL('https://example.com'), configDir: 'fake', logger: loggerTest, version: 'test'}, emitter);
const source = new DeezerTestSource('test', {data: {arl: 'test'}, options}, {localUrl: new URL('https://example.com'), configDir: 'fake', logger: loggerTest, version: 'test'}, emitter);
await source.tryInitialize();
return source;
}
+11 -4
View File
@@ -1,11 +1,18 @@
import { loggerTest } from "@foxxmd/logging";
import { MSCache } from "../../common/Cache.js";
import { getDb, migrateDbSync } from "../../common/database/drizzle/drizzleUtils.js";
import { getDb, migrateDbSync, migrateDb, DbConcrete } from "../../common/database/drizzle/drizzleUtils.js";
import { getPrepopulatedMemoryPGlite } from "./databaseFixtures.js";
import { PGlite } from "@electric-sql/pglite";
export const transientCache = () => new MSCache(loggerTest);
export const transientDb = () => {
const db = getDb(':memory:');
migrateDbSync(db);
let baseDb: PGlite;
export const transientDb = async () => {
if(baseDb === undefined) {
baseDb = await getPrepopulatedMemoryPGlite();
await migrateDb(await getDb(baseDb));
}
const db = getDb((await baseDb.clone()) as Awaited<PGlite>);
return db;
}
@@ -5,6 +5,8 @@ import { PlayNew } from "../../common/database/drizzle/drizzleTypes.js";
import { PlayInputNew } from "../../common/database/drizzle/drizzleTypes.js";
import { ComponentNew } from "../../common/database/drizzle/drizzleTypes.js";
import { ObjectPlayData } from "../../../core/Atomic.js";
import { PGlite } from '@electric-sql/pglite'
import { dataDir } from '@electric-sql/pglite-prepopulatedfs'
export const fixtureCreateComponent = (data: Partial<ComponentNew> = {}): ComponentNew => {
return generateComponentEntity(
@@ -35,4 +37,17 @@ export const fixtureCreateInput = (data: PlayInputNew & { data?: object | false
realData = inputData;
}
return generateInputEntity({...rest, data: realData});
}
export const getPrepopulatedFSPGlite = async (dir: string) => {
return PGlite.create({
dataDir: dir,
loadDataDir: await dataDir()
});
}
export const getPrepopulatedMemoryPGlite = async () => {
return PGlite.create({
loadDataDir: await dataDir()
});
}
+3 -4
View File
@@ -28,7 +28,6 @@ const createYtSource = async (opts?: {
} = opts || {};
const source = new YTMusicSource('test', config, { localUrl: new URL('https://example.com'), configDir: 'fake', logger: loggerTest, version: 'test' }, emitter);
await source.buildDatabase();
source.buildTransformRules();
return source;
}
@@ -140,7 +139,7 @@ describe('Handles temporal inconsistency in history', function () {
const prependedPlays = [newPlay, ...plays];
expect(source.parseRecentAgainstResponse(prependedPlays).plays).length(1);
await sleep(1000);
await sleep(50);
// YT returns outdated history
// should be detected as append since "removed" track in last position from previous history is seen again
@@ -148,12 +147,12 @@ describe('Handles temporal inconsistency in history', function () {
expect(badAppend).to.deep.include({consistent: false, diffType: 'added', plays: []});
expect(badAppend.diffResults[2]).eq('append');
await sleep(500);
await sleep(10);
// contiuned outdated history
expect(source.parseRecentAgainstResponse(plays)).to.deep.include({consistent: true, plays: []});
await sleep(500);
await sleep(10);
// correct, current history is finally returned correctly
const recentHistoryResult = source.parseRecentAgainstResponse(prependedPlays);
+1 -1
View File
@@ -230,7 +230,7 @@ export const splitByFirstRegexFound = <T>(str: any, onNotAStringVal: T, delimsRe
/**
* Returns value if it is a non-empty string or returns default value
* */
export const nonEmptyStringOrDefault = <T>(str: any, defaultVal: T = undefined): string | T => {
export const nonEmptyStringOrDefault = <T = undefined>(str: any, defaultVal: T = undefined): string | T => {
if (str === undefined || str === null || typeof str !== 'string' || str.trim() === '') {
return defaultVal;
}
+4 -1
View File
@@ -33,7 +33,10 @@ export default defineConfig(() => {
console.debug(`[VITE] BASE_URL ENV: ${process.env.BASE_URL} | Base Url String: ${baseUrlStr}`);
return {
server: {
allowedHosts: (true as true)
allowedHosts: (true as true),
watch: {
ignored: ['**/msDb/**', '**/config/**']
}
},
esbuild: {
minifyIdentifiers: false