mirror of
https://github.com/FoxxMD/multi-scrobbler.git
synced 2026-09-03 05:10:00 +03:00
Compare commits
13
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
739ad85ebf | ||
|
|
f3a87911ca | ||
|
|
b9db03f626 | ||
|
|
57f8922866 | ||
|
|
97a01b6be4 | ||
|
|
019133d725 | ||
|
|
303cb4b109 | ||
|
|
479dca660d | ||
|
|
5f0b626b2e | ||
|
|
54a0c74002 | ||
|
|
419eecd5c0 | ||
|
|
63233dc18c | ||
|
|
871882c05b |
@@ -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"
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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
@@ -6,5 +6,6 @@
|
||||
"file": [
|
||||
"./src/backend/tests/setup.ts"
|
||||
],
|
||||
"exit": true
|
||||
"exit": true,
|
||||
"timeout": 2500
|
||||
}
|
||||
|
||||
Vendored
+2
-1
@@ -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
@@ -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'),
|
||||
},
|
||||
});
|
||||
Generated
+36
-8
@@ -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
@@ -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",
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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 => {
|
||||
|
||||
-83
@@ -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`);
|
||||
+83
@@ -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;
|
||||
+499
-239
File diff suppressed because it is too large
Load Diff
@@ -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),
|
||||
|
||||
@@ -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
@@ -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
@@ -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;
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
|
||||
@@ -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();
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
});
|
||||
}
|
||||
@@ -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);
|
||||
|
||||
@@ -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
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user