From 50aabd82913f3d2df5f6ba1226c7213eda1ce0e7 Mon Sep 17 00:00:00 2001 From: ailegion Date: Wed, 9 Sep 2026 11:16:06 +0200 Subject: [PATCH 1/4] refactor: studio database store --- e2e/fixtures/electron-seeded.fixture.ts | 2 + e2e/fixtures/electron.fixture.ts | 2 + src/main/database/index.ts | 12 + src/main/database/migrations.ts | 67 +++++ src/main/database/store.ts | 223 ++++++++++++++ src/main/services/connectors.service.ts | 206 ++++++------- src/main/services/icebergDatalake.service.ts | 119 ++++---- src/main/services/projects.service.ts | 229 ++++++++------- src/main/services/savedQueries.service.ts | 80 +++--- src/main/services/settings.service.ts | 120 ++++---- src/main/utils/fileHelper.ts | 69 +---- src/main/utils/sanitizeBigQueryKeyfile.ts | 22 ++ src/main/utils/setupHelpers.ts | 7 +- src/types/backend.ts | 4 + tests/integration/ipc/settings.ipc.test.ts | 2 + tests/unit/main/database/migrations.test.ts | 78 +++++ tests/unit/main/database/store.test.ts | 272 ++++++++++++++++++ .../main/services/connectors.service.test.ts | 20 +- .../services/icebergDatalake.service.test.ts | 71 +++-- .../main/services/projects.service.test.ts | 49 +++- .../main/services/settings.service.test.ts | 69 +++-- 21 files changed, 1223 insertions(+), 500 deletions(-) create mode 100644 src/main/database/index.ts create mode 100644 src/main/database/migrations.ts create mode 100644 src/main/database/store.ts create mode 100644 src/main/utils/sanitizeBigQueryKeyfile.ts create mode 100644 tests/unit/main/database/migrations.test.ts create mode 100644 tests/unit/main/database/store.test.ts diff --git a/e2e/fixtures/electron-seeded.fixture.ts b/e2e/fixtures/electron-seeded.fixture.ts index 9fefb965..bdf3420b 100644 --- a/e2e/fixtures/electron-seeded.fixture.ts +++ b/e2e/fixtures/electron-seeded.fixture.ts @@ -3,6 +3,7 @@ import { test as base, ElectronApplication, Page } from '@playwright/test'; import * as path from 'path'; import * as fs from 'fs'; import * as os from 'os'; +import { CURRENT_SCHEMA_VERSION } from '../../src/main/database/migrations'; export type TestFixtures = { electronApp: ElectronApplication; @@ -47,6 +48,7 @@ export const test = base.extend({ // Seed database.json with test project and connection const databaseJson = { + schemaVersion: CURRENT_SCHEMA_VERSION, settings: { isSetup: 'true', pythonPath: diff --git a/e2e/fixtures/electron.fixture.ts b/e2e/fixtures/electron.fixture.ts index 795f13e2..e4bfd77b 100644 --- a/e2e/fixtures/electron.fixture.ts +++ b/e2e/fixtures/electron.fixture.ts @@ -13,6 +13,7 @@ import { test as base, ElectronApplication, Page } from '@playwright/test'; import * as path from 'path'; import * as fs from 'fs'; import * as os from 'os'; +import { CURRENT_SCHEMA_VERSION } from '../../src/main/database/migrations'; // Type definitions for our fixtures export type ElectronFixtures = { @@ -81,6 +82,7 @@ export const test = base.extend({ const dbPath = path.join(userData, 'database.json'); const settings = { + schemaVersion: CURRENT_SCHEMA_VERSION, settings: { isSetup: 'true', pythonPath: diff --git a/src/main/database/index.ts b/src/main/database/index.ts new file mode 100644 index 00000000..c3350187 --- /dev/null +++ b/src/main/database/index.ts @@ -0,0 +1,12 @@ +import { DB_FILE } from '../utils/setupHelpers'; +import { DatabaseStore } from './store'; + +// Single shared instance for the whole main process — every service should +// import this rather than constructing its own DatabaseStore, otherwise the +// serialized queue that makes concurrent access safe would be per-instance +// instead of per-file. +const databaseStore = new DatabaseStore(DB_FILE); + +export default databaseStore; +export { DatabaseStore } from './store'; +export { CURRENT_SCHEMA_VERSION } from './migrations'; diff --git a/src/main/database/migrations.ts b/src/main/database/migrations.ts new file mode 100644 index 00000000..f0bd4547 --- /dev/null +++ b/src/main/database/migrations.ts @@ -0,0 +1,67 @@ +import { DataBase } from '../../types/backend'; + +// Bump this whenever a migration is added below. Every database.json ever +// written before this system existed has no `schemaVersion` field at all, +// so it is treated as version 0. +export const CURRENT_SCHEMA_VERSION = 1; + +type RawShape = Record & { schemaVersion?: number }; + +interface Migration { + // The version this migration produces, i.e. it runs for any file whose + // recorded version is lower than this number. + version: number; + migrate: (raw: RawShape) => RawShape; +} + +// Ordered by version, ascending. Each entry transforms the *previous* +// shape into this version's shape. Keep old migrations forever — they are +// what let an install from years ago open successfully today. +const migrations: Migration[] = [ + { + // No structural change: this just stamps every pre-versioning file + // (implicit version 0) so it — and every future file — carries an + // explicit, auditable schema version from here on. + version: 1, + migrate: (raw) => ({ ...raw, schemaVersion: 1 }), + }, +]; + +export function pendingMigrations(fromVersion: number): Migration[] { + return migrations + .filter((m) => m.version > fromVersion) + .sort((a, b) => a.version - b.version); +} + +/** + * Runs every pending migration in order and defensively fills in any field + * still missing afterwards (covers installs older than the oldest + * migration, and callers that hand in a bare `{}`). Never throws on a + * merely-incomplete shape — only a genuinely corrupt/unparseable file + * should reach the store's recovery path instead of this function. + */ +export function migrate(rawInput: unknown): DataBase { + let shape: RawShape = + typeof rawInput === 'object' && rawInput !== null + ? (rawInput as RawShape) + : {}; + + const startVersion = shape.schemaVersion ?? 0; + pendingMigrations(startVersion).forEach((m) => { + shape = m.migrate(shape); + }); + + return { + schemaVersion: CURRENT_SCHEMA_VERSION, + projects: (shape.projects as DataBase['projects']) ?? [], + settings: + (shape.settings as DataBase['settings']) ?? ({} as DataBase['settings']), + selectedProject: shape.selectedProject as DataBase['selectedProject'], + queries: (shape.queries as DataBase['queries']) ?? {}, + savedQueries: shape.savedQueries as DataBase['savedQueries'], + connections: (shape.connections as DataBase['connections']) ?? [], + sources: (shape.sources as DataBase['sources']) ?? [], + recentItems: (shape.recentItems as DataBase['recentItems']) ?? [], + icebergInstances: shape.icebergInstances as DataBase['icebergInstances'], + }; +} diff --git a/src/main/database/store.ts b/src/main/database/store.ts new file mode 100644 index 00000000..3bfcb2e6 --- /dev/null +++ b/src/main/database/store.ts @@ -0,0 +1,223 @@ +import fs from 'fs'; +import path from 'path'; +import { DataBase } from '../../types/backend'; +import { CURRENT_SCHEMA_VERSION, migrate } from './migrations'; + +function defaultDatabase(): DataBase { + return { + schemaVersion: CURRENT_SCHEMA_VERSION, + projects: [], + settings: {} as DataBase['settings'], + queries: {}, + connections: [], + sources: [], + recentItems: [], + }; +} + +/** + * Single point of read/write access to a JSON-file-backed database + * (production: database.json). Replaces the old pattern of independent + * `fs.readFile` + in-memory mutation + `fs.writeFile` call sites, which was + * racy: two concurrent callers could both read the same stale value before + * either wrote back, silently dropping one of their changes. + * + * All access — reads included — goes through a single serialized queue, so + * a `getField` always observes the latest committed write and a + * `updateField`/`transaction` always mutates the true current value rather + * than a snapshot taken before it got in line. + */ +export class DatabaseStore { + private readonly filePath: string; + + private cache: DataBase | null = null; + + private queue: Promise = Promise.resolve(); + + constructor(filePath: string) { + this.filePath = filePath; + } + + /** Drops the in-memory cache so the next access re-reads from disk. */ + invalidateCache(): void { + this.cache = null; + } + + async getField(key: K): Promise { + return this.enqueue(async () => { + const db = await this.ensureLoaded(); + return db[key]; + }); + } + + /** + * Reads multiple fields as one consistent snapshot — use this instead of + * separate getField calls whenever a caller combines two or more fields + * (e.g. joining projects to connections), since two separate getField + * calls are two separate queue turns and a write could land in between. + */ + async getSnapshot(): Promise> { + return this.enqueue(async () => this.ensureLoaded()); + } + + async updateField( + key: K, + updater: (current: DataBase[K]) => DataBase[K], + ): Promise { + return this.enqueue(async () => { + const db = await this.ensureLoaded(); + const next = updater(db[key]); + const updatedDb: DataBase = { ...db, [key]: next }; + await this.persist(updatedDb); + this.cache = updatedDb; + return next; + }); + } + + /** + * For updates that must touch more than one top-level key atomically + * (e.g. updating `selectedProject` and `projects` together) — one lock + * hold, one disk write, instead of two independent updateField calls + * that could otherwise interleave with someone else's write in between. + */ + async transaction( + mutator: (db: DataBase) => { db: DataBase; result: T }, + ): Promise { + return this.enqueue(async () => { + const db = await this.ensureLoaded(); + const { db: nextDb, result } = mutator(db); + await this.persist(nextDb); + this.cache = nextDb; + return result; + }); + } + + private enqueue(operation: () => Promise): Promise { + const result = this.queue.then(operation, operation); + // Keep the chain alive regardless of this operation's outcome so one + // failed read/write never permanently jams every later caller. + this.queue = result.then( + () => undefined, + () => undefined, + ); + return result; + } + + private async ensureLoaded(): Promise { + if (!this.cache) { + this.cache = await this.loadFromDisk(); + } + return this.cache; + } + + private async loadFromDisk(): Promise { + let raw: string; + try { + raw = await fs.promises.readFile(this.filePath, 'utf8'); + } catch (err: unknown) { + if ((err as { code?: string })?.code === 'ENOENT') { + // Don't write anything on a mere read of a not-yet-existing + // profile — the file is created naturally by the first real + // write (updateField/transaction), same as before this store + // existed. Avoids surprising disk writes on read-only access. + return defaultDatabase(); + } + throw err; + } + + let parsed: unknown; + try { + parsed = JSON.parse(raw); + } catch (parseError) { + return this.recoverFromBackup(parseError as Error); + } + + const onDiskVersion = + (parsed as { schemaVersion?: number })?.schemaVersion ?? 0; + const needsMigration = onDiskVersion < CURRENT_SCHEMA_VERSION; + if (needsMigration) { + await this.backup(raw, onDiskVersion); + } + const migrated = migrate(parsed); + if (needsMigration) { + await this.persist(migrated); + } + return migrated; + } + + private backupPath(fromVersion: number): string { + return `${this.filePath}.v${fromVersion}.bak-${Date.now()}`; + } + + private async backup(raw: string, fromVersion: number): Promise { + await fs.promises.writeFile(this.backupPath(fromVersion), raw, 'utf8'); + } + + /** + * A file that fails to JSON.parse is corruption, not "no file" — treating + * it as an empty install (the old loadDatabaseFile behavior) would + * silently destroy every project/connection/setting the user has. Try + * the most recent pre-migration backup instead; only give up and throw + * if none parses either. + */ + private async recoverFromBackup(parseError: Error): Promise { + const dir = path.dirname(this.filePath); + const base = path.basename(this.filePath); + let entries: string[] = []; + try { + entries = await fs.promises.readdir(dir); + } catch { + // fall through to the "no backup" error below + } + + const backups = entries + .filter((name) => name.startsWith(`${base}.v`) && name.includes('.bak-')) + .sort() + .reverse(); + + // eslint-disable-next-line no-restricted-syntax + for (const backupName of backups) { + try { + // eslint-disable-next-line no-await-in-loop + const backupRaw = await fs.promises.readFile( + path.join(dir, backupName), + 'utf8', + ); + const parsed = JSON.parse(backupRaw); + // eslint-disable-next-line no-console + console.error( + `[DatabaseStore] ${this.filePath} is corrupted (${parseError.message}); recovered from backup ${backupName}`, + ); + return migrate(parsed); + } catch { + // Try the next-oldest backup. + } + } + + throw new Error( + `${this.filePath} is corrupted and no valid backup could be recovered: ${parseError.message}`, + ); + } + + /** Atomic write: stage in a temp file, fsync, then rename over the real + * path. A crash or failed write mid-way leaves the previous file intact + * instead of a truncated/corrupted one. */ + private async persist(data: DataBase): Promise { + const dir = path.dirname(this.filePath); + // Don't assume some other init step already created the directory — + // the store should be self-sufficient on a completely fresh profile. + await fs.promises.mkdir(dir, { recursive: true }); + const tmpPath = path.join( + dir, + `.${path.basename(this.filePath)}.tmp-${process.pid}-${Date.now()}-${Math.random().toString(36).slice(2)}`, + ); + const handle = await fs.promises.open(tmpPath, 'w'); + try { + await handle.writeFile(JSON.stringify(data, null, 2), 'utf8'); + await handle.sync(); + } finally { + await handle.close(); + } + await fs.promises.rename(tmpPath, this.filePath); + } +} diff --git a/src/main/services/connectors.service.ts b/src/main/services/connectors.service.ts index c9eaed12..00718061 100644 --- a/src/main/services/connectors.service.ts +++ b/src/main/services/connectors.service.ts @@ -24,7 +24,8 @@ import { SnowflakeConnection, SQLiteConnection, } from '../../types/backend'; -import { loadDatabaseFile, updateDatabase } from '../utils/fileHelper'; +import databaseStore from '../database'; +import { sanitizeBigQueryKeyfile } from '../utils/sanitizeBigQueryKeyfile'; import { ProjectsService } from './index'; import MainDatabaseService from './mainDatabase.service'; import { ConfigureConnectionBody, UpdateConnectionBody } from '../../types/ipc'; @@ -133,8 +134,7 @@ export default class ConnectorsService { static async loadConnections( includeDataLake: boolean = false, ): Promise { - const db = await loadDatabaseFile(); - const connections = db.connections ?? []; + const connections = (await databaseStore.getField('connections')) ?? []; // Filter out ducklake connections by default if (includeDataLake) { @@ -343,45 +343,33 @@ export default class ConnectorsService { connection: ConnectionInput, allowReservedNames: boolean = false, ): Promise { - const connections = await this.loadConnections(true); // Include all connections including ducklake - - // Validate connection name with optional allowReservedNames flag - const nameValidation = this.validateConnectionName( - connection.name, - connections, - undefined, - allowReservedNames, - ); - - if (!nameValidation.isValid) { - throw new Error(nameValidation.message); - } - const connectionId = uuidV4(); const newConnection: ConnectionModel = { id: connectionId, connection, }; - await updateDatabase<'connections'>('connections', [ - ...connections, - newConnection, - ]); + + await databaseStore.updateField('connections', (current) => { + const connections = current ?? []; + // Validate connection name with optional allowReservedNames flag + const nameValidation = this.validateConnectionName( + connection.name, + connections, + undefined, + allowReservedNames, + ); + if (!nameValidation.isValid) { + throw new Error(nameValidation.message); + } + const next = [...connections, newConnection]; + next.forEach((conn) => sanitizeBigQueryKeyfile(conn)); + return next; + }); + return connectionId; } static async saveNewConnection(connection: ConnectionInput): Promise { - const connections = await this.loadConnections(true); // Include all connections including ducklake - - // Validate connection name - const nameValidation = this.validateConnectionName( - connection.name, - connections, - ); - - if (!nameValidation.isValid) { - throw new Error(nameValidation.message); - } - const connectionId = uuidV4(); // For ducklake connections, store S3 credentials securely @@ -420,15 +408,26 @@ export default class ConnectorsService { id: connectionId, connection, }; - await updateDatabase<'connections'>('connections', [ - ...connections, - newConnection, - ]); + + await databaseStore.updateField('connections', (current) => { + const connections = current ?? []; + const nameValidation = this.validateConnectionName( + connection.name, + connections, + ); + if (!nameValidation.isValid) { + throw new Error(nameValidation.message); + } + const next = [...connections, newConnection]; + next.forEach((conn) => sanitizeBigQueryKeyfile(conn)); + return next; + }); + return connectionId; } static async getProjectById(projectId: string): Promise { - const { projects } = await loadDatabaseFile(); + const projects = await databaseStore.getField('projects'); return projects.find((p) => p.id === projectId); } @@ -621,29 +620,30 @@ export default class ConnectorsService { }: UpdateConnectionBody): Promise { await this.validateConnection(connection.connection); - const connections = await this.loadConnections(true); // Include all connections including ducklake - - // Validate connection name (exclude current connection from uniqueness check) - const nameValidation = this.validateConnectionName( - connection.connection.name, - connections, - connection.id, - ); - - if (!nameValidation.isValid) { - throw new Error(nameValidation.message); - } - - const connectionIndex = connections.findIndex( - (c) => c.id === connection.id, - ); + await databaseStore.updateField('connections', (current) => { + const connections = current ?? []; + // Validate connection name (exclude current connection from uniqueness check) + const nameValidation = this.validateConnectionName( + connection.connection.name, + connections, + connection.id, + ); + if (!nameValidation.isValid) { + throw new Error(nameValidation.message); + } - if (connectionIndex === -1) { - throw new Error('Connection not found'); - } + const connectionIndex = connections.findIndex( + (c) => c.id === connection.id, + ); + if (connectionIndex === -1) { + throw new Error('Connection not found'); + } - connections[connectionIndex] = connection; - await updateDatabase<'connections'>('connections', connections); + const updated = [...connections]; + updated[connectionIndex] = connection; + updated.forEach((conn) => sanitizeBigQueryKeyfile(conn)); + return updated; + }); if (connection.connection.type === 'sqlite') return; @@ -744,12 +744,10 @@ export default class ConnectorsService { await NotebooksService.archiveConnectionNotebooks(connectionToDelete.id); // Remove the connection from the database - const updatedConnections = connections.filter( - (connection) => connection.id !== connectionId, + await databaseStore.updateField('connections', (current) => + (current ?? []).filter((c) => c.id !== connectionId), ); - await updateDatabase<'connections'>('connections', updatedConnections); - // Only clean up AI chats after the connection deletion is persisted. try { await MainDatabaseService.deleteConversationsByConnection(connectionId); @@ -1793,29 +1791,25 @@ export default class ConnectorsService { } static async loadCloudConnections(): Promise { - const db = await loadDatabaseFile(); - return db.sources ?? []; + const sources = await databaseStore.getField('sources'); + return sources ?? []; } static async saveCloudConnection(connection: CloudConnection): Promise { - const db = await loadDatabaseFile(); - const sources = db.sources ?? []; - - const existingIndex = sources.findIndex((c) => c.id === connection.id); - - if (existingIndex >= 0) { - sources[existingIndex] = connection; - } else { - sources.push(connection); - } - - await updateDatabase<'sources'>('sources', sources); + await databaseStore.updateField('sources', (current) => { + const sources = current ?? []; + const existingIndex = sources.findIndex((c) => c.id === connection.id); + if (existingIndex >= 0) { + const updated = [...sources]; + updated[existingIndex] = connection; + return updated; + } + return [...sources, connection]; + }); } static async deleteCloudConnection(id: string): Promise { - const db = await loadDatabaseFile(); - const sources = db.sources ?? []; - + const sources = (await databaseStore.getField('sources')) ?? []; const connectionToDelete = sources.find((c) => c.id === id); if (connectionToDelete) { // Clean up cloud connection-specific credentials from secure storage @@ -1832,9 +1826,9 @@ export default class ConnectorsService { } } - const filteredSources = sources.filter((c) => c.id !== id); - - await updateDatabase<'sources'>('sources', filteredSources); + await databaseStore.updateField('sources', (current) => + (current ?? []).filter((c) => c.id !== id), + ); } static async getCloudConnectionById( @@ -1850,9 +1844,8 @@ export default class ConnectorsService { static async loadRecentItems(): Promise { try { - const db = await loadDatabaseFile(); - const items = db.recentItems ?? []; - return items.sort( + const items = (await databaseStore.getField('recentItems')) ?? []; + return [...items].sort( (a, b) => new Date(b.accessedAt).getTime() - new Date(a.accessedAt).getTime(), ); @@ -1864,32 +1857,21 @@ export default class ConnectorsService { static async addRecentItem( item: Omit, ): Promise { - const db = await loadDatabaseFile(); - const items = db.recentItems ?? []; - - const existingIndex = items.findIndex((i) => i.id === item.id); - - if (existingIndex >= 0) { - items.splice(existingIndex, 1); - } - - items.unshift({ ...item, accessedAt: new Date() }); - - const recentItems = items.slice(0, 50); - await updateDatabase<'recentItems'>('recentItems', recentItems); + await databaseStore.updateField('recentItems', (current) => { + const items = (current ?? []).filter((i) => i.id !== item.id); + const next = [{ ...item, accessedAt: new Date() }, ...items]; + return next.slice(0, 50); + }); } static async removeRecentItem(id: string): Promise { - const db = await loadDatabaseFile(); - const items = db.recentItems ?? []; - - const filteredItems = items.filter((i) => i.id !== id); - - await updateDatabase<'recentItems'>('recentItems', filteredItems); + await databaseStore.updateField('recentItems', (current) => + (current ?? []).filter((i) => i.id !== id), + ); } static async clearRecentItems(): Promise { - await updateDatabase<'recentItems'>('recentItems', []); + await databaseStore.updateField('recentItems', () => []); } /** @@ -2110,19 +2092,19 @@ export default class ConnectorsService { connectionId: string, query: string, ): Promise { - const db = await loadDatabaseFile(); - const queries = db.queries ?? {}; // Use a connection-specific key prefix to distinguish from project queries - queries[`connection:${connectionId}`] = query; - await updateDatabase('queries', queries); + await databaseStore.updateField('queries', (current) => ({ + ...(current ?? {}), + [`connection:${connectionId}`]: query, + })); } /** * Get the saved query for a specific connection */ static async getConnectionQuery(connectionId: string): Promise { - const db = await loadDatabaseFile(); - return db.queries?.[`connection:${connectionId}`] ?? ''; + const queries = await databaseStore.getField('queries'); + return queries?.[`connection:${connectionId}`] ?? ''; } /** diff --git a/src/main/services/icebergDatalake.service.ts b/src/main/services/icebergDatalake.service.ts index e41fa50b..7357fe4c 100644 --- a/src/main/services/icebergDatalake.service.ts +++ b/src/main/services/icebergDatalake.service.ts @@ -14,7 +14,7 @@ import { pathToFileURL } from 'url'; import { spawn } from 'child_process'; import { app } from 'electron'; -import { loadDatabaseFile, updateDatabase } from '../utils/fileHelper'; +import databaseStore from '../database'; import secureStorage from './secureStorage.service'; import SettingsService from './settings.service'; @@ -131,14 +131,23 @@ export class IcebergDatalakeService { // Private: persistence helpers // ───────────────────────────────────────────── + // Old persisted records used catalogType 'file'; normalize to 'sqlite' on + // every read. Any subsequent write of the normalized array also fixes the + // value at rest, so this doubles as a lazy migration for that field. + private static normalizeInstances( + instances: IcebergInstanceConfig[], + ): IcebergInstanceConfig[] { + return instances.map((instance) => { + const persisted = instance as unknown as { catalogType: string }; + if (persisted.catalogType !== 'file') return instance; + return { ...instance, catalogType: 'sqlite' }; + }); + } + private static async readInstances(): Promise { try { - const db = await loadDatabaseFile(); - return (db.icebergInstances ?? []).map((instance) => { - const persisted = instance as unknown as { catalogType: string }; - if (persisted.catalogType !== 'file') return instance; - return { ...instance, catalogType: 'sqlite' }; - }); + const instances = await databaseStore.getField('icebergInstances'); + return IcebergDatalakeService.normalizeInstances(instances ?? []); } catch (error) { // eslint-disable-next-line no-console console.error('[IcebergDatalakeService] readInstances error:', error); @@ -146,12 +155,6 @@ export class IcebergDatalakeService { } } - private static async writeInstances( - instances: IcebergInstanceConfig[], - ): Promise { - await updateDatabase('icebergInstances', instances); - } - // ───────────────────────────────────────────── // Private: Python bridge helpers // ───────────────────────────────────────────── @@ -467,8 +470,8 @@ export class IcebergDatalakeService { if (instance.storageConnectionId) { try { - const db = await loadDatabaseFile(); - const conn: CloudConnection | undefined = (db.sources ?? []).find( + const sources = await databaseStore.getField('sources'); + const conn: CloudConnection | undefined = (sources ?? []).find( (s) => s.id === instance.storageConnectionId, ); if (conn) { @@ -570,8 +573,8 @@ export class IcebergDatalakeService { if (!config.databaseConnectionId) { throw new Error('ICEBERG_REQUIRED_FIELD: databaseConnectionId'); } - const db = await loadDatabaseFile(); - const model = (db.connections ?? []).find( + const connections = await databaseStore.getField('connections'); + const model = (connections ?? []).find( (item) => item.id === config.databaseConnectionId, ); if (!model || model.connection.type !== 'postgres') { @@ -823,9 +826,10 @@ export class IcebergDatalakeService { updatedAt: now, }; - const instances = await IcebergDatalakeService.readInstances(); - instances.push(newInstance); - await IcebergDatalakeService.writeInstances(instances); + await databaseStore.updateField('icebergInstances', (current) => [ + ...IcebergDatalakeService.normalizeInstances(current ?? []), + newInstance, + ]); return newInstance; } catch (error) { @@ -840,30 +844,34 @@ export class IcebergDatalakeService { data: UpdateIcebergInstanceDTO, ): Promise { try { - const instances = await IcebergDatalakeService.readInstances(); - const idx = instances.findIndex((i) => i.id === id); - if (idx < 0) throw new Error(`Iceberg instance not found: ${id}`); + const existing = (await IcebergDatalakeService.readInstances()).find( + (i) => i.id === id, + ); + if (!existing) throw new Error(`Iceberg instance not found: ${id}`); - const updatedConfig = { - ...instances[idx], - ...data, - }; + const updatedConfig = { ...existing, ...data }; IcebergDatalakeService.validateCatalogWarehousePair(updatedConfig); IcebergDatalakeService.validateCatalogAuthentication(updatedConfig); // Handle access token update + let { catalogAccessTokenKey } = existing; if (data.accessToken) { - const key = - instances[idx].catalogAccessTokenKey ?? `iceberg-catalog-token-${id}`; - await secureStorage.setCredential(key, data.accessToken); - instances[idx].catalogAccessTokenKey = key; + catalogAccessTokenKey = + catalogAccessTokenKey ?? `iceberg-catalog-token-${id}`; + await secureStorage.setCredential( + catalogAccessTokenKey, + data.accessToken, + ); } + let { oauthClientSecretKey } = existing; if (data.oauthClientSecret) { - const key = - instances[idx].oauthClientSecretKey ?? `iceberg-oauth-secret-${id}`; - await secureStorage.setCredential(key, data.oauthClientSecret); - instances[idx].oauthClientSecretKey = key; + oauthClientSecretKey = + oauthClientSecretKey ?? `iceberg-oauth-secret-${id}`; + await secureStorage.setCredential( + oauthClientSecretKey, + data.oauthClientSecret, + ); } // Strip raw secrets before persisting @@ -875,15 +883,28 @@ export class IcebergDatalakeService { } = data; /* eslint-enable @typescript-eslint/no-unused-vars */ - instances[idx] = { - ...instances[idx], + const changes = { ...rest, id, + catalogAccessTokenKey, + oauthClientSecretKey, updatedAt: new Date().toISOString(), }; - await IcebergDatalakeService.writeInstances(instances); - return instances[idx]; + let finalInstance: IcebergInstanceConfig | undefined; + await databaseStore.updateField('icebergInstances', (current) => { + const instances = IcebergDatalakeService.normalizeInstances( + current ?? [], + ); + const idx = instances.findIndex((i) => i.id === id); + if (idx < 0) throw new Error(`Iceberg instance not found: ${id}`); + const next = [...instances]; + next[idx] = { ...next[idx], ...changes }; + finalInstance = next[idx]; + return next; + }); + + return finalInstance as IcebergInstanceConfig; } catch (error) { // eslint-disable-next-line no-console console.error('[IcebergDatalakeService] updateInstance error:', error); @@ -893,8 +914,9 @@ export class IcebergDatalakeService { static async deleteInstance(id: string): Promise { try { - const instances = await IcebergDatalakeService.readInstances(); - const instance = instances.find((i) => i.id === id); + const instance = (await IcebergDatalakeService.readInstances()).find( + (i) => i.id === id, + ); if (!instance) throw new Error(`Iceberg instance not found: ${id}`); if (instance.catalogAccessTokenKey) { @@ -920,8 +942,11 @@ export class IcebergDatalakeService { } } - const updated = instances.filter((i) => i.id !== id); - await IcebergDatalakeService.writeInstances(updated); + await databaseStore.updateField('icebergInstances', (current) => + IcebergDatalakeService.normalizeInstances(current ?? []).filter( + (i) => i.id !== id, + ), + ); } catch (error) { // eslint-disable-next-line no-console console.error('[IcebergDatalakeService] deleteInstance error:', error); @@ -1146,8 +1171,8 @@ export class IcebergDatalakeService { } private static async resolveCloudStorageConnection(connectionId: string) { - const db = await loadDatabaseFile(); - const connection = (db.sources ?? []).find( + const sources = await databaseStore.getField('sources'); + const connection = (sources ?? []).find( (source) => source.id === connectionId, ); if (!connection) throw new Error('ICEBERG_CLOUD_CONNECTION_NOT_FOUND'); @@ -1238,7 +1263,6 @@ export class IcebergDatalakeService { // Settings record the last successful installation for diagnostics, but // do not prove the currently selected Python still has every required // extra. Verify once per app session before trusting it. - const settings = await SettingsService.loadSettings(); // Check via Python bridge (runs pip only if the runtime profile is incomplete) const checkResult = (await IcebergDatalakeService.runBridge({ @@ -1248,7 +1272,6 @@ export class IcebergDatalakeService { if (checkResult.installed) { const version = checkResult.version as string | undefined; await SettingsService.saveSettings({ - ...settings, icebergInstalled: true, icebergVersion: version, }); @@ -1283,9 +1306,7 @@ export class IcebergDatalakeService { if (verifyResult.installed) { const version = verifyResult.version as string | undefined; - const currentSettings = await SettingsService.loadSettings(); await SettingsService.saveSettings({ - ...currentSettings, icebergInstalled: true, icebergVersion: version, }); diff --git a/src/main/services/projects.service.ts b/src/main/services/projects.service.ts index ef468515..ecbed7bb 100644 --- a/src/main/services/projects.service.ts +++ b/src/main/services/projects.service.ts @@ -9,6 +9,7 @@ import * as tar from 'tar'; import { BigQueryConnection, DatabricksConnection, + DataBase, DuckDBConnection, KineticaConnection, PostgresConnection, @@ -26,12 +27,12 @@ import { deleteDirectory, deleteItem, getDirectoryStructure, - loadDatabaseFile, readFileContent, saveFileContent, searchInFiles, - updateDatabase, } from '../utils/fileHelper'; +import databaseStore from '../database'; +import { sanitizeBigQueryKeyfile } from '../utils/sanitizeBigQueryKeyfile'; import SettingsService from './settings.service'; import { BigQueryExtractor, @@ -51,10 +52,20 @@ import { } from '../utils/pipelineEnvVars'; export default class ProjectsService { + // Mirrors loadProjects()'s join against raw store state — used inside + // transactions where the raw (un-joined) DataBase is all that's available. + private static joinProjects( + db: Pick, + ): Project[] { + return db.projects.map((p) => ({ + ...p, + connection: db.connections.find((c) => c.id === p.connectionId) + ?.connection, + })); + } + static async loadProjects(): Promise { - const db = await loadDatabaseFile(); - const { connections } = db; - const { projects } = db; + const { projects, connections } = await databaseStore.getSnapshot(); return projects.map((project) => ({ ...project, @@ -91,35 +102,24 @@ export default class ProjectsService { } static async getSelectedProject(): Promise { - const db = await loadDatabaseFile(); - const selected = db.selectedProject; + const selected = await databaseStore.getField('selectedProject'); try { const project = await this.getProject(selected?.id); if (!project) { - await updateDatabase<'selectedProject'>('selectedProject', undefined); + await databaseStore.updateField('selectedProject', () => undefined); return undefined; } return project; } catch (err) { // If loading the project or its configuration fails, clear selection - await updateDatabase<'selectedProject'>('selectedProject', undefined); + await databaseStore.updateField('selectedProject', () => undefined); return undefined; } } static async saveProjects(projects: Project[]) { - // Patch: For all projects, if the connection is bigquery and keyfile is a JSON string, store only the key name - for (const project of projects) { - if ( - project.connection && - project.connection.type === 'bigquery' && - project.connection.keyfile && - project.connection.keyfile.startsWith('{') - ) { - project.connection.keyfile = `db-bigquery-${project.connection.name}`; - } - } - await updateDatabase<'projects'>('projects', projects); + projects.forEach((project) => sanitizeBigQueryKeyfile(project)); + await databaseStore.updateField('projects', () => projects); } static async addProject( @@ -140,15 +140,7 @@ export default class ProjectsService { createTemplateFolders, }; - // Patch: If the project has a bigquery connection, store only the key name - if (project.connection && project.connection.type === 'bigquery') { - if ( - project.connection.keyfile && - project.connection.keyfile.startsWith('{') - ) { - project.connection.keyfile = `db-bigquery-${project.connection.name}`; - } - } + sanitizeBigQueryKeyfile(project); await this.copyDbtTemplateFiles(project.path, project.name); // Always copy main.conf template (rosetta requires it), but profiles.yml is excluded from template await this.copyRosettaMainConf(project.path); @@ -231,15 +223,7 @@ export default class ProjectsService { connectionId, }; - // Patch: If the project has a bigquery connection, store only the key name - if (project.connection && project.connection.type === 'bigquery') { - if ( - project.connection.keyfile && - project.connection.keyfile.startsWith('{') - ) { - project.connection.keyfile = `db-bigquery-${project.connection.name}`; - } - } + sanitizeBigQueryKeyfile(project); const rosettaPath = path.join(projectPath, 'rosetta'); @@ -647,37 +631,59 @@ export default class ProjectsService { } static async updateProject(project: Project) { - const projects = await this.loadProjects(); - const index = projects.findIndex((p) => p.id === project.id); - if (index === -1) return null; - - // Check if connectionId is changing - const oldConnectionId = projects[index].connectionId; - const newConnectionId = project.connectionId; - const connectionChanged = oldConnectionId !== newConnectionId; + // `projects` (the full list) and `selectedProject` (a denormalized copy + // of one entry) must move together — updating them as two separate + // writes could leave selectedProject pointing at data inconsistent with + // the list if anything failed or interleaved in between. + const { updatedProject, connectionChanged, projects } = + await databaseStore.transaction<{ + updatedProject: Project | null; + connectionChanged: boolean; + projects: Project[] | null; + }>((db) => { + // Mirrors loadProjects()'s join so the persisted array keeps the + // same (denormalized but harmless — always re-derived on next + // read) shape it always has. + const joinedProjects = this.joinProjects(db); + const index = joinedProjects.findIndex((p) => p.id === project.id); + if (index === -1) { + return { + db, + result: { + updatedProject: null, + connectionChanged: false, + projects: null, + }, + }; + } - const updatedProject = { ...projects[index], ...project }; + const oldConnectionId = joinedProjects[index].connectionId; + const newConnectionId = project.connectionId; + const merged = sanitizeBigQueryKeyfile({ + ...joinedProjects[index], + ...project, + }); + joinedProjects[index] = merged; - // Patch: If the project has a bigquery connection, store only the key name - if ( - updatedProject.connection && - updatedProject.connection.type === 'bigquery' - ) { - if ( - updatedProject.connection.keyfile && - updatedProject.connection.keyfile.startsWith('{') - ) { - updatedProject.connection.keyfile = `db-bigquery-${updatedProject.connection.name}`; - } - } + return { + db: { + ...db, + projects: joinedProjects, + selectedProject: merged, + }, + result: { + updatedProject: merged, + connectionChanged: oldConnectionId !== newConnectionId, + projects: joinedProjects, + }, + }; + }); - projects[index] = updatedProject; - await updateDatabase<'selectedProject'>('selectedProject', updatedProject); - await this.saveProjects(projects); + if (!updatedProject) return null; // Only regenerate config files if the connection changed // This is a full regeneration because it's switching to a different connection - if (connectionChanged && newConnectionId) { + if (connectionChanged && project.connectionId) { await ConnectorsService.loadConfigurations(project.id); } @@ -687,34 +693,40 @@ export default class ProjectsService { static async deleteProject(id: string) { const projects = await this.loadProjects(); const projectToDelete = projects.find((p) => p.id === id); - if (projectToDelete) { - if (projectToDelete.path) { - deleteDirectory(projectToDelete.path); - } - const selectedProject = await this.getSelectedProject(); - if (selectedProject) { - if (selectedProject.id === id) { - await updateDatabase('selectedProject', undefined); - } - } + if (!projectToDelete) return false; - const filteredProjects = projects.filter((p) => p.id !== id); - await this.saveProjects(filteredProjects); + if (projectToDelete.path) { + deleteDirectory(projectToDelete.path); + } - // Only clean up AI chats after the project deletion is persisted. - try { - const projectIdNum = parseInt(id, 10); - if (!Number.isNaN(projectIdNum)) { - await MainDatabaseService.deleteConversationsByProject(projectIdNum); - } - } catch (error) { - // eslint-disable-next-line no-console - console.error('[ProjectsService] Failed to clean up AI chats:', error); - } + // Both the list and the (possibly now-dangling) selection move together + // in one write — see updateProject for why that matters. + await databaseStore.transaction((db) => { + const remaining = this.joinProjects(db).filter((p) => p.id !== id); + remaining.forEach((p) => sanitizeBigQueryKeyfile(p)); + return { + db: { + ...db, + projects: remaining, + selectedProject: + db.selectedProject?.id === id ? undefined : db.selectedProject, + }, + result: undefined, + }; + }); - return true; + // Only clean up AI chats after the project deletion is persisted. + try { + const projectIdNum = parseInt(id, 10); + if (!Number.isNaN(projectIdNum)) { + await MainDatabaseService.deleteConversationsByProject(projectIdNum); + } + } catch (error) { + // eslint-disable-next-line no-console + console.error('[ProjectsService] Failed to clean up AI chats:', error); } - return false; + + return true; } // Removes a project from Studio's known-projects list without touching @@ -723,18 +735,23 @@ export default class ProjectsService { static async removeProjectFromList(id: string) { const projects = await this.loadProjects(); const projectToRemove = projects.find((p) => p.id === id); - if (projectToRemove) { - const selectedProject = await this.getSelectedProject(); - if (selectedProject && selectedProject.id === id) { - await updateDatabase('selectedProject', undefined); - } - - const filteredProjects = projects.filter((p) => p.id !== id); - await this.saveProjects(filteredProjects); + if (!projectToRemove) return false; + + await databaseStore.transaction((db) => { + const remaining = this.joinProjects(db).filter((p) => p.id !== id); + remaining.forEach((p) => sanitizeBigQueryKeyfile(p)); + return { + db: { + ...db, + projects: remaining, + selectedProject: + db.selectedProject?.id === id ? undefined : db.selectedProject, + }, + result: undefined, + }; + }); - return true; - } - return false; + return true; } static async getProjectPath(name: string) { @@ -1018,7 +1035,7 @@ export default class ProjectsService { static async selectProject({ projectId }: { projectId: string }) { const project = await this.getProject(projectId); - await updateDatabase<'selectedProject'>('selectedProject', project); + await databaseStore.updateField('selectedProject', () => project); } static async extractPgSchema(connection: PostgresConnection) { @@ -1233,15 +1250,15 @@ export default class ProjectsService { projectId: string; query: string; }): Promise { - const db = await loadDatabaseFile(); - const queries = db.queries ?? {}; - queries[projectId] = query; - await updateDatabase('queries', queries); + await databaseStore.updateField('queries', (current) => ({ + ...(current ?? {}), + [projectId]: query, + })); } static async getQuery(project: Project): Promise { - const db = await loadDatabaseFile(); - return db.queries?.[project.id] ?? ''; + const queries = await databaseStore.getField('queries'); + return queries?.[project.id] ?? ''; } static async extractSchemaFromModelYaml(project: Project): Promise { diff --git a/src/main/services/savedQueries.service.ts b/src/main/services/savedQueries.service.ts index 894501c1..3ad65873 100644 --- a/src/main/services/savedQueries.service.ts +++ b/src/main/services/savedQueries.service.ts @@ -1,5 +1,5 @@ import { v4 as uuidv4 } from 'uuid'; -import { loadDatabaseFile, updateDatabase } from '../utils/fileHelper'; +import databaseStore from '../database'; import { SavedQuery } from '../../types/backend'; export class SavedQueriesService { @@ -7,9 +7,8 @@ export class SavedQueriesService { * List saved queries for a specific connection */ static async list(connectionId: string): Promise { - const db = await loadDatabaseFile(); - const savedQueries = db.savedQueries || {}; - return savedQueries[connectionId] || []; + const savedQueries = await databaseStore.getField('savedQueries'); + return savedQueries?.[connectionId] || []; } /** @@ -20,10 +19,6 @@ export class SavedQueriesService { name: string, query: string, ): Promise { - const db = await loadDatabaseFile(); - const savedQueries = db.savedQueries || {}; - const connectionQueries = savedQueries[connectionId] || []; - const newQuery: SavedQuery = { id: uuidv4(), name, @@ -33,11 +28,13 @@ export class SavedQueriesService { updatedAt: new Date().toISOString(), }; - const updatedConnectionQueries = [...connectionQueries, newQuery]; - - await updateDatabase('savedQueries', { - ...savedQueries, - [connectionId]: updatedConnectionQueries, + await databaseStore.updateField('savedQueries', (current) => { + const savedQueries = current || {}; + const connectionQueries = savedQueries[connectionId] || []; + return { + ...savedQueries, + [connectionId]: [...connectionQueries, newQuery], + }; }); return newQuery; @@ -51,47 +48,46 @@ export class SavedQueriesService { queryId: string, updates: Partial>, ): Promise { - const db = await loadDatabaseFile(); - const savedQueries = db.savedQueries || {}; - const connectionQueries = savedQueries[connectionId] || []; + let updatedQuery: SavedQuery | undefined; - const queryIndex = connectionQueries.findIndex((q) => q.id === queryId); - if (queryIndex === -1) { - throw new Error(`Saved query not found: ${queryId}`); - } + await databaseStore.updateField('savedQueries', (current) => { + const savedQueries = current || {}; + const connectionQueries = savedQueries[connectionId] || []; - const updatedQuery: SavedQuery = { - ...connectionQueries[queryIndex], - ...updates, - updatedAt: new Date().toISOString(), - }; + const queryIndex = connectionQueries.findIndex((q) => q.id === queryId); + if (queryIndex === -1) { + throw new Error(`Saved query not found: ${queryId}`); + } - const updatedConnectionQueries = [...connectionQueries]; - updatedConnectionQueries[queryIndex] = updatedQuery; + updatedQuery = { + ...connectionQueries[queryIndex], + ...updates, + updatedAt: new Date().toISOString(), + }; - await updateDatabase('savedQueries', { - ...savedQueries, - [connectionId]: updatedConnectionQueries, + const updatedConnectionQueries = [...connectionQueries]; + updatedConnectionQueries[queryIndex] = updatedQuery; + + return { + ...savedQueries, + [connectionId]: updatedConnectionQueries, + }; }); - return updatedQuery; + return updatedQuery as SavedQuery; } /** * Delete a saved query */ static async delete(connectionId: string, queryId: string): Promise { - const db = await loadDatabaseFile(); - const savedQueries = db.savedQueries || {}; - const connectionQueries = savedQueries[connectionId] || []; - - const updatedConnectionQueries = connectionQueries.filter( - (q) => q.id !== queryId, - ); - - await updateDatabase('savedQueries', { - ...savedQueries, - [connectionId]: updatedConnectionQueries, + await databaseStore.updateField('savedQueries', (current) => { + const savedQueries = current || {}; + const connectionQueries = savedQueries[connectionId] || []; + return { + ...savedQueries, + [connectionId]: connectionQueries.filter((q) => q.id !== queryId), + }; }); } } diff --git a/src/main/services/settings.service.ts b/src/main/services/settings.service.ts index 5999c268..585b5553 100644 --- a/src/main/services/settings.service.ts +++ b/src/main/services/settings.service.ts @@ -7,11 +7,8 @@ import type { Session } from 'electron'; import os from 'os'; import AdmZip from 'adm-zip'; import * as tar from 'tar'; -import { - loadDatabaseFile, - loadDefaultSettings, - updateDatabase, -} from '../utils/fileHelper'; +import { loadDefaultSettings } from '../utils/fileHelper'; +import databaseStore from '../database'; import { CliUpdateResponseType, SettingsType, @@ -68,18 +65,16 @@ export default class SettingsService { private static factoryResetPromise: Promise | null = null; static async loadSettings(): Promise { - const dataBase = await loadDatabaseFile(); + const settings = await databaseStore.getField('settings'); const defaultSettings = loadDefaultSettings(); return { ...defaultSettings, - ...dataBase.settings, + ...settings, projectsDirectory: - dataBase.settings.projectsDirectory || - defaultSettings.projectsDirectory, + settings.projectsDirectory || defaultSettings.projectsDirectory, sampleRosettaMainConf: - dataBase.settings.sampleRosettaMainConf || - defaultSettings.sampleRosettaMainConf, + settings.sampleRosettaMainConf || defaultSettings.sampleRosettaMainConf, }; } @@ -144,8 +139,16 @@ export default class SettingsService { return DuckDBBootstrap.diagnose(); } - static async saveSettings(settings: SettingsType) { - await updateDatabase<'settings'>('settings', settings); + // Takes a partial update, merged against the current settings inside the + // store's atomic write — not a full object read earlier by the caller. + // Two install/uninstall flows running concurrently (e.g. installing + // Python while installing Rosetta) now each only touch their own fields + // instead of one clobbering the other's change with a stale full copy. + static async saveSettings(partialSettings: Partial) { + await databaseStore.updateField('settings', (current) => ({ + ...current, + ...partialSettings, + })); } static async getDbtExePath(): Promise { @@ -218,7 +221,6 @@ export default class SettingsService { static async updateRosetta() { if (process.env.E2E_TESTING === 'true') { - const settings = await this.loadSettings(); const dummyName = process.platform === 'win32' ? 'dummy-rosetta.exe' : 'dummy-rosetta'; const dummyPath = path.join(os.tmpdir(), dummyName); @@ -228,9 +230,10 @@ export default class SettingsService { fs.chmodSync(dummyPath, 0o755); } - settings.rosettaVersion = '0.0.0-test'; - settings.rosettaPath = dummyPath; - await this.saveSettings(settings); + await this.saveSettings({ + rosettaVersion: '0.0.0-test', + rosettaPath: dummyPath, + }); return { binaryPath: dummyPath, @@ -330,9 +333,10 @@ export default class SettingsService { }), ); - settings.rosettaVersion = version; - settings.rosettaPath = binaryPath; - await this.saveSettings(settings); + await this.saveSettings({ + rosettaVersion: version, + rosettaPath: binaryPath, + }); await fs.remove(zipPath); @@ -346,8 +350,6 @@ export default class SettingsService { private static async performPythonInstall(version: string) { if (process.env.E2E_TESTING === 'true') { - const settings = await this.loadSettings(); - // Create a dummy venv structure const venvPath = path.join(os.tmpdir(), 'dummy-venv'); const binDir = @@ -363,10 +365,11 @@ export default class SettingsService { await fs.ensureFile(dummyBinaryPath); await fs.chmod(dummyBinaryPath, 0o755); - settings.pythonVersion = version; - settings.pythonPath = dummyBinaryPath; - settings.pythonBinary = dummyBinaryPath; - await this.saveSettings(settings); + await this.saveSettings({ + pythonVersion: version, + pythonPath: dummyBinaryPath, + pythonBinary: dummyBinaryPath, + }); return { binaryPath: dummyBinaryPath, @@ -459,21 +462,24 @@ export default class SettingsService { // fool every "is Python configured" check elsewhere into skipping // auto-install/repair and silently falling back to the unmanaged, // non-isolated base interpreter. - await this.clearManagedVenvDependents(settings); + const clearedDependents = await this.clearManagedVenvDependents(); await fs.remove(venvDir); - settings.pythonVersion = ''; - settings.pythonPath = ''; - settings.pythonBinary = ''; - await this.saveSettings(settings); + await this.saveSettings({ + ...clearedDependents, + pythonVersion: '', + pythonPath: '', + pythonBinary: '', + }); const cliAdapter = new CliAdapter(); await cliAdapter.runCommandWithoutStreaming( `cd "${userDataPath}" && "${binaryPath}" -m venv venv`, ); - settings.pythonVersion = version; - settings.pythonPath = venvPythonPath; - settings.pythonBinary = binaryPath; - await this.saveSettings(settings); + await this.saveSettings({ + pythonVersion: version, + pythonPath: venvPythonPath, + pythonBinary: binaryPath, + }); await fs.remove(archivePath); return { @@ -556,9 +562,9 @@ export default class SettingsService { } } - private static async clearManagedVenvDependents( - settings: SettingsType, - ): Promise { + private static async clearManagedVenvDependents(): Promise< + Partial + > { try { const { FlowfileService } = await import('./flowfile.service'); await FlowfileService.stop(); @@ -568,27 +574,23 @@ export default class SettingsService { // dbt-core (v1 and v2), Flowfile, and sqlglot are all pip-installed into // the same managed venv, so removing it also removes their executables. - // eslint-disable-next-line no-param-reassign - settings.dbtPath = ''; - // eslint-disable-next-line no-param-reassign - settings.dbtVersion = ''; - // eslint-disable-next-line no-param-reassign - settings.flowfileVersion = ''; + return { dbtPath: '', dbtVersion: '', flowfileVersion: '' }; } static async uninstallPython(): Promise { - const settings = await this.loadSettings(); const userDataPath = app.getPath('userData'); - await this.clearManagedVenvDependents(settings); + const clearedDependents = await this.clearManagedVenvDependents(); await fs.remove(path.join(userDataPath, 'python')); await fs.remove(path.join(userDataPath, 'venv')); - settings.pythonVersion = ''; - settings.pythonPath = ''; - settings.pythonBinary = ''; - await this.saveSettings(settings); + await this.saveSettings({ + ...clearedDependents, + pythonVersion: '', + pythonPath: '', + pythonBinary: '', + }); } static async resetFactorySettings(session: Session): Promise { @@ -608,12 +610,13 @@ export default class SettingsService { let teardownStarted = false; try { - const dataBase = await loadDatabaseFile(); + const projects = await databaseStore.getField('projects'); + const settings = await databaseStore.getField('settings'); const projectPaths = await this.resolveProjectPaths( - (dataBase.projects ?? []).map((project) => project.path), + (projects ?? []).map((project) => project.path), ); const managedRosettaPath = await this.resolveManagedRosettaPath( - dataBase.settings?.rosettaPath, + settings?.rosettaPath, ); stage = 'stopping active resources'; @@ -1080,9 +1083,10 @@ export default class SettingsService { ); // Update settings - settings.rosettaVersion = version; - settings.rosettaPath = binaryPath; - await this.saveSettings(settings); + await this.saveSettings({ + rosettaVersion: version, + rosettaPath: binaryPath, + }); // Clean up download file await fs.remove(zipPath); @@ -1110,9 +1114,7 @@ export default class SettingsService { await fs.remove(rosettaRoot); } - settings.rosettaVersion = ''; - settings.rosettaPath = ''; - await this.saveSettings(settings); + await this.saveSettings({ rosettaVersion: '', rosettaPath: '' }); } private static getRosettaDownloadUrl(release: any): string { diff --git a/src/main/utils/fileHelper.ts b/src/main/utils/fileHelper.ts index 884a0c25..39bf5f45 100644 --- a/src/main/utils/fileHelper.ts +++ b/src/main/utils/fileHelper.ts @@ -4,13 +4,12 @@ import { app, dialog } from 'electron'; import archiver from 'archiver'; import os from 'os'; import { - DataBase, FileNode, FileSearchMatch, FileSearchResult, SettingsType, } from '../../types/backend'; -import { DATA_DIR, DB_FILE } from './setupHelpers'; +import { DATA_DIR } from './setupHelpers'; export const getDirectoryStructure = (dirPath: string): FileNode => { const result: FileNode = { @@ -252,72 +251,6 @@ export const loadDefaultSettings = (): SettingsType => { }; }; -export const loadDatabaseFile = async (): Promise => { - try { - const data = await fs.promises.readFile(DB_FILE, 'utf8'); - const parsed = JSON.parse(data) as Partial; - return { - ...parsed, - projects: parsed.projects ?? [], - settings: parsed.settings ?? loadDefaultSettings(), - queries: parsed.queries ?? {}, - connections: parsed.connections ?? [], - sources: parsed.sources ?? [], - recentItems: parsed.recentItems ?? [], - }; - } catch (error) { - return { - projects: [], - settings: loadDefaultSettings(), - queries: {}, - connections: [], - sources: [], - recentItems: [], - }; - } -}; - -// Simple async mutex to prevent concurrent read-modify-write races -let dbLockPromise: Promise = Promise.resolve(); - -export const updateDatabase = async ( - key: K, - value: DataBase[K], -) => { - // Chain writes so they execute sequentially - const previousLock = dbLockPromise; - let releaseLock: () => void; - dbLockPromise = new Promise((resolve) => { - releaseLock = resolve; - }); - - try { - await previousLock; - - // Patch: For connections array, ensure BigQuery keyfile is only the key name - if (key === 'connections' && Array.isArray(value)) { - value.forEach((conn) => { - if ( - conn && - typeof conn === 'object' && - 'connection' in conn && - conn.connection && - conn.connection.type === 'bigquery' && - conn.connection.keyfile && - conn.connection.keyfile.startsWith('{') - ) { - conn.connection.keyfile = `db-bigquery-${conn.connection.name}`; - } - }); - } - const data = await loadDatabaseFile(); - data[key] = value; - await saveFileContent(DB_FILE, JSON.stringify(data, null, 2)); - } finally { - releaseLock!(); - } -}; - export const createNewFolder = (parentPath: string, folderName: string) => { const folderPath = path.join(parentPath, folderName); diff --git a/src/main/utils/sanitizeBigQueryKeyfile.ts b/src/main/utils/sanitizeBigQueryKeyfile.ts new file mode 100644 index 00000000..f81efa63 --- /dev/null +++ b/src/main/utils/sanitizeBigQueryKeyfile.ts @@ -0,0 +1,22 @@ +/** + * Raw BigQuery service-account JSON must never be persisted to database.json + * — only the secure-storage key name it was materialized under. This was + * previously duplicated ad hoc in four separate call sites; centralized here + * so there's exactly one place to fix if the key-naming scheme ever changes. + */ +export function sanitizeBigQueryKeyfile< + T extends { + connection?: { type?: string; keyfile?: string; name?: string }; + }, +>(item: T): T { + const { connection } = item; + if ( + connection && + connection.type === 'bigquery' && + connection.keyfile && + connection.keyfile.startsWith('{') + ) { + connection.keyfile = `db-bigquery-${connection.name}`; + } + return item; +} diff --git a/src/main/utils/setupHelpers.ts b/src/main/utils/setupHelpers.ts index 2d8252d6..34662a80 100644 --- a/src/main/utils/setupHelpers.ts +++ b/src/main/utils/setupHelpers.ts @@ -12,9 +12,10 @@ export const initializeDataStorage = async () => { fs.mkdirSync(DATA_DIR, { recursive: true }); } - if (!fs.existsSync(DB_FILE)) { - fs.writeFileSync(DB_FILE, JSON.stringify({ projects: [] }, null, 2)); - } + // database.json itself no longer needs eager seeding here: DatabaseStore + // creates it (with a valid, versioned shape) automatically the first time + // any service reads or writes through it, which happens moments after + // this function returns. // Initialize main database for AI and future features try { diff --git a/src/types/backend.ts b/src/types/backend.ts index 993d44c8..c3f68b15 100644 --- a/src/types/backend.ts +++ b/src/types/backend.ts @@ -375,6 +375,10 @@ export interface SavedQuery { } export type DataBase = { + // Absent on any database.json written before the migration system existed + // (schema version 0). Never written as undefined once loaded through + // DatabaseStore — always stamped to CURRENT_SCHEMA_VERSION on save. + schemaVersion?: number; projects: Project[]; settings: SettingsType; selectedProject?: Project; diff --git a/tests/integration/ipc/settings.ipc.test.ts b/tests/integration/ipc/settings.ipc.test.ts index d5441d1e..ade0a92d 100644 --- a/tests/integration/ipc/settings.ipc.test.ts +++ b/tests/integration/ipc/settings.ipc.test.ts @@ -4,6 +4,7 @@ import * as fs from 'fs'; // Import handlers after mocks import registerSettingsHandlers from '../../../src/main/ipcHandlers/settings.ipcHandlers'; +import { CURRENT_SCHEMA_VERSION } from '../../../src/main/database/migrations'; // Define path constants const TEST_DIR_NAME = 'dbt-studio-settings-ipc-test'; @@ -96,6 +97,7 @@ describe('Settings IPC Integration', () => { // Create database.json file with default settings const dbPath = path.join(MOCK_USER_DATA, 'database.json'); const defaultDb = { + schemaVersion: CURRENT_SCHEMA_VERSION, projects: [], settings: { dbtVersion: '1.0.0', diff --git a/tests/unit/main/database/migrations.test.ts b/tests/unit/main/database/migrations.test.ts new file mode 100644 index 00000000..d515f696 --- /dev/null +++ b/tests/unit/main/database/migrations.test.ts @@ -0,0 +1,78 @@ +import { + CURRENT_SCHEMA_VERSION, + migrate, + pendingMigrations, +} from '../../../../src/main/database/migrations'; + +describe('migrations', () => { + describe('migrate', () => { + it('fills in every field with a safe default from an empty object', () => { + const result = migrate({}); + + expect(result).toEqual({ + schemaVersion: CURRENT_SCHEMA_VERSION, + projects: [], + settings: {}, + selectedProject: undefined, + queries: {}, + savedQueries: undefined, + connections: [], + sources: [], + recentItems: [], + icebergInstances: undefined, + }); + }); + + it('treats non-object input (e.g. a corrupt-but-parseable value) the same as empty', () => { + expect(migrate(null).schemaVersion).toBe(CURRENT_SCHEMA_VERSION); + expect(migrate(42).projects).toEqual([]); + expect(migrate('oops').connections).toEqual([]); + }); + + it('stamps a pre-versioning (implicit version 0) shape without touching its data', () => { + const legacy = { + projects: [{ id: '1', name: 'demo' }], + connections: [{ id: 'c1' }], + }; + + const result = migrate(legacy); + + expect(result.schemaVersion).toBe(CURRENT_SCHEMA_VERSION); + expect(result.projects).toEqual(legacy.projects); + expect(result.connections).toEqual(legacy.connections); + // Fields absent from the legacy file still get their defaults. + expect(result.recentItems).toEqual([]); + expect(result.sources).toEqual([]); + }); + + it('is a no-op on data that is already at the current version', () => { + const current = { + schemaVersion: CURRENT_SCHEMA_VERSION, + projects: [{ id: '1' }], + settings: { isSetup: 'true' }, + queries: {}, + connections: [], + sources: [], + recentItems: [], + }; + + const result = migrate(current); + + expect(result).toEqual(current); + }); + }); + + describe('pendingMigrations', () => { + it('returns every migration above the given version, in ascending order', () => { + const pending = pendingMigrations(0); + expect(pending.map((m) => m.version)).toEqual( + [...pending.map((m) => m.version)].sort((a, b) => a - b), + ); + expect(pending.every((m) => m.version > 0)).toBe(true); + }); + + it('returns nothing once already at the current version', () => { + expect(pendingMigrations(CURRENT_SCHEMA_VERSION)).toEqual([]); + }); + }); +}); diff --git a/tests/unit/main/database/store.test.ts b/tests/unit/main/database/store.test.ts new file mode 100644 index 00000000..809397be --- /dev/null +++ b/tests/unit/main/database/store.test.ts @@ -0,0 +1,272 @@ +import fs from 'fs'; +import os from 'os'; +import path from 'path'; +import { DatabaseStore } from '../../../../src/main/database/store'; +import { CURRENT_SCHEMA_VERSION } from '../../../../src/main/database/migrations'; + +describe('DatabaseStore', () => { + let dir: string; + let dbPath: string; + + beforeEach(() => { + dir = fs.mkdtempSync(path.join(os.tmpdir(), 'db-store-test-')); + dbPath = path.join(dir, 'database.json'); + }); + + afterEach(() => { + fs.rmSync(dir, { recursive: true, force: true }); + }); + + describe('initialization', () => { + it('returns in-memory defaults without touching disk when the file does not exist yet', async () => { + const store = new DatabaseStore(dbPath); + + const projects = await store.getField('projects'); + + expect(projects).toEqual([]); + expect(fs.existsSync(dbPath)).toBe(false); + }); + + it('creates the file (and schema-stamps it) on the first write', async () => { + const store = new DatabaseStore(dbPath); + + await store.updateField('projects', () => [{ id: 'p1' } as never]); + + expect(fs.existsSync(dbPath)).toBe(true); + const onDisk = JSON.parse(fs.readFileSync(dbPath, 'utf8')); + expect(onDisk.schemaVersion).toBe(CURRENT_SCHEMA_VERSION); + expect(onDisk.projects).toEqual([{ id: 'p1' }]); + }); + + it('creates its parent directory on the first write if it does not exist yet', async () => { + const freshDir = path.join(dir, 'nested', 'profile-dir'); + const freshDbPath = path.join(freshDir, 'database.json'); + const store = new DatabaseStore(freshDbPath); + + await store.updateField('projects', () => [{ id: 'p1' } as never]); + + expect(fs.existsSync(freshDbPath)).toBe(true); + }); + + it('loads an existing up-to-date file as-is', async () => { + fs.writeFileSync( + dbPath, + JSON.stringify({ + schemaVersion: CURRENT_SCHEMA_VERSION, + projects: [{ id: 'p1', name: 'existing' }], + settings: {}, + queries: {}, + connections: [], + sources: [], + recentItems: [], + }), + ); + + const store = new DatabaseStore(dbPath); + const projects = await store.getField('projects'); + + expect(projects).toEqual([{ id: 'p1', name: 'existing' }]); + }); + }); + + describe('migration on load', () => { + it('migrates a pre-versioning file, defaults missing fields, and backs up the original', async () => { + const legacyContent = JSON.stringify({ + projects: [{ id: 'old-project' }], + }); + fs.writeFileSync(dbPath, legacyContent); + + const store = new DatabaseStore(dbPath); + const projects = await store.getField('projects'); + const recentItems = await store.getField('recentItems'); + + expect(projects).toEqual([{ id: 'old-project' }]); + expect(recentItems).toEqual([]); + + const onDisk = JSON.parse(fs.readFileSync(dbPath, 'utf8')); + expect(onDisk.schemaVersion).toBe(CURRENT_SCHEMA_VERSION); + + const backupFiles = fs + .readdirSync(dir) + .filter((name) => name.includes('.v0.bak-')); + expect(backupFiles).toHaveLength(1); + expect(fs.readFileSync(path.join(dir, backupFiles[0]), 'utf8')).toBe( + legacyContent, + ); + }); + }); + + describe('corruption recovery', () => { + it('recovers from the most recent backup when the live file fails to parse', async () => { + fs.writeFileSync(dbPath, '{ this is not valid json'); + const goodBackup = { + schemaVersion: CURRENT_SCHEMA_VERSION, + projects: [{ id: 'recovered-project' }], + settings: {}, + queries: {}, + connections: [], + sources: [], + recentItems: [], + }; + fs.writeFileSync(`${dbPath}.v1.bak-1000`, JSON.stringify(goodBackup)); + + const store = new DatabaseStore(dbPath); + const projects = await store.getField('projects'); + + expect(projects).toEqual([{ id: 'recovered-project' }]); + }); + + it('picks the most recent backup when several exist', async () => { + fs.writeFileSync(dbPath, '{ still not valid json'); + fs.writeFileSync( + `${dbPath}.v1.bak-1000`, + JSON.stringify({ projects: [{ id: 'older-backup' }] }), + ); + fs.writeFileSync( + `${dbPath}.v1.bak-2000`, + JSON.stringify({ projects: [{ id: 'newer-backup' }] }), + ); + + const store = new DatabaseStore(dbPath); + const projects = await store.getField('projects'); + + expect(projects).toEqual([{ id: 'newer-backup' }]); + }); + + it('throws instead of silently wiping data when no valid backup exists', async () => { + fs.writeFileSync(dbPath, '{ not valid json at all'); + + const store = new DatabaseStore(dbPath); + + await expect(store.getField('projects')).rejects.toThrow(/corrupted/i); + }); + }); + + describe('updateField', () => { + it('persists the new value and later reads observe it', async () => { + const store = new DatabaseStore(dbPath); + + await store.updateField('recentItems', () => [{ id: 'a' } as never]); + const recentItems = await store.getField('recentItems'); + + expect(recentItems).toEqual([{ id: 'a' }]); + const onDisk = JSON.parse(fs.readFileSync(dbPath, 'utf8')); + expect(onDisk.recentItems).toEqual([{ id: 'a' }]); + }); + + it('does not lose updates when many concurrent callers touch the same key', async () => { + const store = new DatabaseStore(dbPath); + await store.getField('projects'); // warm the in-memory cache first + + const concurrentWrites = 25; + await Promise.all( + Array.from({ length: concurrentWrites }, (_, i) => + store.updateField('recentItems', (current) => [ + ...current, + { id: `item-${i}` } as never, + ]), + ), + ); + + const recentItems = await store.getField('recentItems'); + expect(recentItems).toHaveLength(concurrentWrites); + const ids = new Set(recentItems.map((item: any) => item.id)); + expect(ids.size).toBe(concurrentWrites); + }); + + it('applies concurrent updates to different keys without clobbering each other', async () => { + const store = new DatabaseStore(dbPath); + await store.getField('projects'); + + await Promise.all([ + store.updateField('projects', () => [{ id: 'p1' } as never]), + store.updateField('sources', () => [{ id: 's1' } as never]), + ]); + + expect(await store.getField('projects')).toEqual([{ id: 'p1' }]); + expect(await store.getField('sources')).toEqual([{ id: 's1' }]); + }); + + it('leaves the on-disk file untouched if the write fails partway through', async () => { + const store = new DatabaseStore(dbPath); + await store.updateField('recentItems', () => [{ id: 'before' } as never]); + const beforeContent = fs.readFileSync(dbPath, 'utf8'); + + const renameSpy = jest + .spyOn(fs.promises, 'rename') + .mockRejectedValueOnce(new Error('simulated disk failure')); + + await expect( + store.updateField('recentItems', () => [{ id: 'after' } as never]), + ).rejects.toThrow('simulated disk failure'); + + renameSpy.mockRestore(); + + expect(fs.readFileSync(dbPath, 'utf8')).toBe(beforeContent); + expect(await store.getField('recentItems')).toEqual([{ id: 'before' }]); + }); + }); + + describe('transaction', () => { + it('applies multiple key changes as a single atomic write', async () => { + const store = new DatabaseStore(dbPath); + await store.getField('projects'); + + const writeSpy = jest.spyOn(fs.promises, 'rename'); + const callsBefore = writeSpy.mock.calls.length; + + const result = await store.transaction((db) => ({ + db: { + ...db, + projects: [{ id: 'p1' } as never], + selectedProject: { id: 'p1' } as never, + }, + result: 'done', + })); + + expect(result).toBe('done'); + expect(writeSpy.mock.calls.length).toBe(callsBefore + 1); + expect(await store.getField('projects')).toEqual([{ id: 'p1' }]); + expect(await store.getField('selectedProject')).toEqual({ id: 'p1' }); + }); + }); + + describe('getSnapshot', () => { + it('reflects the latest write and combines fields consistently', async () => { + const store = new DatabaseStore(dbPath); + await store.updateField('projects', () => [{ id: 'p1' } as never]); + await store.updateField('connections', () => [{ id: 'c1' } as never]); + + const snapshot = await store.getSnapshot(); + + expect(snapshot.projects).toEqual([{ id: 'p1' }]); + expect(snapshot.connections).toEqual([{ id: 'c1' }]); + }); + }); + + describe('invalidateCache', () => { + it('forces the next access to re-read from disk', async () => { + const store = new DatabaseStore(dbPath); + await store.updateField('recentItems', () => [{ id: 'a' } as never]); + + // Simulate an external process rewriting the file (e.g. factory reset). + fs.writeFileSync( + dbPath, + JSON.stringify({ + schemaVersion: CURRENT_SCHEMA_VERSION, + projects: [], + settings: {}, + queries: {}, + connections: [], + sources: [], + recentItems: [{ id: 'external-write' }], + }), + ); + + store.invalidateCache(); + expect(await store.getField('recentItems')).toEqual([ + { id: 'external-write' }, + ]); + }); + }); +}); diff --git a/tests/unit/main/services/connectors.service.test.ts b/tests/unit/main/services/connectors.service.test.ts index a1da23de..c7bd5d5b 100644 --- a/tests/unit/main/services/connectors.service.test.ts +++ b/tests/unit/main/services/connectors.service.test.ts @@ -11,11 +11,21 @@ jest.mock('openai', () => ({ OpenAI: jest.fn(), })); -jest.mock('../../../../src/main/utils/fileHelper', () => ({ - loadDatabaseFile: jest - .fn() - .mockResolvedValue({ connections: [], projects: [] }), - updateDatabase: jest.fn(), +jest.mock('../../../../src/main/database', () => ({ + __esModule: true, + default: { + getField: jest.fn((key: string) => + Promise.resolve( + key === 'connections' || key === 'projects' ? [] : undefined, + ), + ), + updateField: jest.fn( + (key: string, updater: (current: unknown) => unknown) => + Promise.resolve( + updater(key === 'connections' || key === 'projects' ? [] : undefined), + ), + ), + }, })); jest.mock('../../../../src/main/services/index', () => ({ diff --git a/tests/unit/main/services/icebergDatalake.service.test.ts b/tests/unit/main/services/icebergDatalake.service.test.ts index 962d8574..b89ad087 100644 --- a/tests/unit/main/services/icebergDatalake.service.test.ts +++ b/tests/unit/main/services/icebergDatalake.service.test.ts @@ -1,13 +1,13 @@ import { IcebergDatalakeService } from '../../../../src/main/services/icebergDatalake.service'; import secureStorage from '../../../../src/main/services/secureStorage.service'; -import { - loadDatabaseFile, - updateDatabase, -} from '../../../../src/main/utils/fileHelper'; - -jest.mock('../../../../src/main/utils/fileHelper', () => ({ - loadDatabaseFile: jest.fn(), - updateDatabase: jest.fn(), +import databaseStore from '../../../../src/main/database'; + +jest.mock('../../../../src/main/database', () => ({ + __esModule: true, + default: { + getField: jest.fn(), + updateField: jest.fn(), + }, })); jest.mock('../../../../src/main/services/secureStorage.service', () => ({ @@ -24,15 +24,34 @@ jest.mock('../../../../src/main/services/settings.service', () => ({ default: { loadSettings: jest.fn() }, })); -const mockedLoadDatabase = loadDatabaseFile as jest.Mock; -const mockedUpdateDatabase = updateDatabase as jest.Mock; +const mockedGetField = databaseStore.getField as jest.Mock; +const mockedUpdateField = databaseStore.updateField as jest.Mock; const mockedSecureStorage = secureStorage as jest.Mocked; describe('IcebergDatalakeService compatibility and secret persistence', () => { + // Lightweight fake backing store: getField/updateField mutate this in + // place, mirroring the real DatabaseStore's contract closely enough for + // these tests (which care about what ends up persisted, not the storage + // mechanism itself). + let icebergInstances: any[]; + beforeEach(() => { jest.clearAllMocks(); - mockedLoadDatabase.mockResolvedValue({ icebergInstances: [] }); - mockedUpdateDatabase.mockResolvedValue(undefined); + icebergInstances = []; + mockedGetField.mockImplementation((key: string) => + Promise.resolve( + key === 'icebergInstances' ? icebergInstances : undefined, + ), + ); + mockedUpdateField.mockImplementation( + (key: string, updater: (current: unknown) => unknown) => { + if (key === 'icebergInstances') { + icebergInstances = updater(icebergInstances) as any[]; + return Promise.resolve(icebergInstances); + } + return Promise.resolve(updater(undefined)); + }, + ); mockedSecureStorage.setCredential.mockResolvedValue(undefined); }); @@ -117,7 +136,7 @@ describe('IcebergDatalakeService compatibility and secret persistence', () => { `iceberg-oauth-secret-${created.id}`, 'top-secret', ); - const persisted = mockedUpdateDatabase.mock.calls[0][1][0]; + const persisted = icebergInstances[0]; expect(persisted.oauthClientSecret).toBeUndefined(); expect(JSON.stringify(persisted)).not.toContain('top-secret'); expect(persisted.oauthClientSecretKey).toBe( @@ -135,7 +154,7 @@ describe('IcebergDatalakeService compatibility and secret persistence', () => { storageType: 'server-managed', }); - const persisted = mockedUpdateDatabase.mock.calls[0][1][0]; + const persisted = icebergInstances[0]; expect(persisted).toMatchObject({ id: created.id, catalogType: 'lakekeeper', @@ -427,19 +446,17 @@ describe('IcebergDatalakeService compatibility and secret persistence', () => { }); it('passes sanitized nested namespaces to the bridge', async () => { - mockedLoadDatabase.mockResolvedValue({ - icebergInstances: [ - { - id: 'instance-1', - name: 'test', - catalogType: 'sqlite', - storageType: 'local', - localPath: '/tmp/warehouse', - createdAt: 'now', - updatedAt: 'now', - }, - ], - }); + icebergInstances = [ + { + id: 'instance-1', + name: 'test', + catalogType: 'sqlite', + storageType: 'local', + localPath: '/tmp/warehouse', + createdAt: 'now', + updatedAt: 'now', + }, + ]; const runBridgeSpy = jest .spyOn(IcebergDatalakeService as any, 'runBridge') .mockResolvedValue({ ok: true, namespace: ['a', 'b'] }); diff --git a/tests/unit/main/services/projects.service.test.ts b/tests/unit/main/services/projects.service.test.ts index b97e27e2..42c93c5b 100644 --- a/tests/unit/main/services/projects.service.test.ts +++ b/tests/unit/main/services/projects.service.test.ts @@ -2,12 +2,7 @@ jest.mock('openai', () => ({ OpenAI: jest.fn(), })); -const loadDatabaseFile = jest.fn(); -const updateDatabase = jest.fn(); - jest.mock('../../../../src/main/utils/fileHelper', () => ({ - loadDatabaseFile: (...args: any[]) => loadDatabaseFile(...args), - updateDatabase: (...args: any[]) => updateDatabase(...args), createNewFile: jest.fn(), createNewFolder: jest.fn(), copyPath: jest.fn(), @@ -19,6 +14,16 @@ jest.mock('../../../../src/main/utils/fileHelper', () => ({ saveFileContent: jest.fn(), })); +jest.mock('../../../../src/main/database', () => ({ + __esModule: true, + default: { + getField: jest.fn(), + updateField: jest.fn(), + transaction: jest.fn(), + getSnapshot: jest.fn(), + }, +})); + jest.mock('../../../../src/main/services/settings.service', () => ({ __esModule: true, default: { @@ -57,15 +62,34 @@ jest.mock('../../../../src/main/extractor', () => ({ })); import ProjectsService from '../../../../src/main/services/projects.service'; +import databaseStore from '../../../../src/main/database'; + +const mockedGetSnapshot = databaseStore.getSnapshot as jest.Mock; +const mockedTransaction = databaseStore.transaction as jest.Mock; describe('ProjectsService (main)', () => { + // Lightweight fake backing store: getSnapshot reads it, transaction + // mutates it in place — mirrors DatabaseStore's contract closely enough + // for these tests, which care about persisted state, not the storage + // mechanism. + let fakeDb: { connections: unknown[]; projects: unknown[] }; + beforeEach(() => { jest.clearAllMocks(); + fakeDb = { connections: [], projects: [] }; + mockedGetSnapshot.mockImplementation(() => Promise.resolve(fakeDb)); + mockedTransaction.mockImplementation( + (mutator: (db: typeof fakeDb) => { db: typeof fakeDb; result: unknown }) => { + const { db, result } = mutator(fakeDb); + fakeDb = db; + return Promise.resolve(result); + }, + ); }); describe('loadProjects', () => { it('maps connection config onto projects', async () => { - loadDatabaseFile.mockResolvedValue({ + fakeDb = { connections: [ { id: 'c1', @@ -82,7 +106,7 @@ describe('ProjectsService (main)', () => { isExtracted: false, }, ], - }); + }; const result = await ProjectsService.loadProjects(); expect(result).toHaveLength(1); @@ -94,7 +118,7 @@ describe('ProjectsService (main)', () => { it('updates lastOpenedAt and returns configured project when ConnectorsService.loadConfigurations succeeds', async () => { const nowSpy = jest.spyOn(Date, 'now').mockReturnValue(123); - loadDatabaseFile.mockResolvedValue({ + fakeDb = { connections: [], projects: [ { @@ -105,7 +129,7 @@ describe('ProjectsService (main)', () => { isExtracted: false, }, ], - }); + }; parseProjectConnectionFiles.mockResolvedValue({ rosettaConnection: { dialect: 'duckdb' }, @@ -114,8 +138,7 @@ describe('ProjectsService (main)', () => { const result = await ProjectsService.getProject('p1'); - expect(updateDatabase).toHaveBeenCalledWith( - 'projects', + expect(fakeDb.projects).toEqual( expect.arrayContaining([ expect.objectContaining({ id: 'p1', lastOpenedAt: 123 }), ]), @@ -133,7 +156,7 @@ describe('ProjectsService (main)', () => { }); it('falls back to raw project when ConnectorsService.loadConfigurations throws', async () => { - loadDatabaseFile.mockResolvedValue({ + fakeDb = { connections: [], projects: [ { @@ -144,7 +167,7 @@ describe('ProjectsService (main)', () => { isExtracted: false, }, ], - }); + }; parseProjectConnectionFiles.mockImplementation(async () => { throw new Error('Missing connection'); diff --git a/tests/unit/main/services/settings.service.test.ts b/tests/unit/main/services/settings.service.test.ts index 3f367e38..a4346ae8 100644 --- a/tests/unit/main/services/settings.service.test.ts +++ b/tests/unit/main/services/settings.service.test.ts @@ -7,6 +7,7 @@ import MainDatabaseService from '../../../../src/main/services/mainDatabase.serv import { TaskManagerService } from '../../../../src/main/services/taskManager.service'; import { FlowfileService } from '../../../../src/main/services/flowfile.service'; import * as fileHelper from '../../../../src/main/utils/fileHelper'; +import databaseStore from '../../../../src/main/database'; jest.mock('openai', () => ({ OpenAI: jest.fn(), @@ -77,12 +78,18 @@ jest.mock('../../../../src/main/services/mainDatabase.service', () => ({ })); jest.mock('../../../../src/main/utils/fileHelper', () => ({ - loadDatabaseFile: jest.fn(), - updateDatabase: jest.fn(), loadDefaultSettings: jest.fn(), deleteDirectory: jest.fn(), })); +jest.mock('../../../../src/main/database', () => ({ + __esModule: true, + default: { + getField: jest.fn(), + updateField: jest.fn(), + }, +})); + jest.mock('../../../../src/main/utils/setupHelpers', () => ({ DB_FILE: '/tmp/db.json', initializeDataStorage: jest.fn(), @@ -94,10 +101,18 @@ jest.mock('../../../../src/main/adapters', () => ({ })), })); -const loadDatabaseFile = fileHelper.loadDatabaseFile as jest.Mock; -const updateDatabase = fileHelper.updateDatabase as jest.Mock; +const mockedGetField = databaseStore.getField as jest.Mock; +const mockedUpdateField = databaseStore.updateField as jest.Mock; const loadDefaultSettings = fileHelper.loadDefaultSettings as jest.Mock; +// Stubs databaseStore.getField to serve fixed values per key, mirroring +// what a real DatabaseStore would hand back for the given fields. +const stubDb = (fields: Record) => { + mockedGetField.mockImplementation((key: string) => + Promise.resolve(fields[key]), + ); +}; + describe('SettingsService (main)', () => { beforeEach(() => { jest.clearAllMocks(); @@ -119,26 +134,48 @@ describe('SettingsService (main)', () => { describe('loadSettings', () => { it('returns existing settings when present in database', async () => { - loadDatabaseFile.mockResolvedValue({ - settings: { pythonPath: '/x/python' }, - }); + stubDb({ settings: { pythonPath: '/x/python' } }); loadDefaultSettings.mockReturnValue({ pythonPath: '/default/python' }); await expect(SettingsService.loadSettings()).resolves.toEqual({ pythonPath: '/x/python', }); expect(loadDefaultSettings).toHaveBeenCalled(); - expect(updateDatabase).not.toHaveBeenCalled(); + expect(mockedUpdateField).not.toHaveBeenCalled(); }); it('returns default settings when persisted settings are missing', async () => { - loadDatabaseFile.mockResolvedValue({ settings: {} }); + stubDb({ settings: {} }); loadDefaultSettings.mockReturnValue({ pythonPath: '/default/python' }); await expect(SettingsService.loadSettings()).resolves.toEqual({ pythonPath: '/default/python', }); - expect(updateDatabase).not.toHaveBeenCalled(); + expect(mockedUpdateField).not.toHaveBeenCalled(); + }); + }); + + describe('saveSettings', () => { + it('merges the partial update onto the current settings instead of replacing them', async () => { + await SettingsService.saveSettings({ rosettaVersion: '2.0.0' }); + + expect(mockedUpdateField).toHaveBeenCalledWith( + 'settings', + expect.any(Function), + ); + const updater = mockedUpdateField.mock.calls[0][1]; + const current = { + rosettaVersion: '1.0.0', + pythonPath: '/usr/bin/python3', + }; + + // This is the actual regression guard: a concurrent caller changing + // an unrelated field (pythonPath) must survive this save untouched — + // the old full-object saveSettings would have silently dropped it. + expect(updater(current)).toEqual({ + rosettaVersion: '2.0.0', + pythonPath: '/usr/bin/python3', + }); }); }); @@ -151,7 +188,7 @@ describe('SettingsService (main)', () => { it('stops resources, clears browser data and credentials, deletes owned state, and schedules restart', async () => { jest.useFakeTimers(); - loadDatabaseFile.mockResolvedValue({ + stubDb({ projects: [{ path: '/tmp/dbt-studio-project' }], settings: { rosettaPath: @@ -189,7 +226,7 @@ describe('SettingsService (main)', () => { }); it('does not restart when browser cleanup fails', async () => { - loadDatabaseFile.mockResolvedValue({ projects: [], settings: {} }); + stubDb({ projects: [], settings: {} }); const session = makeSession(); session.clearCache.mockRejectedValue(new Error('/secret/cache/path')); @@ -206,7 +243,7 @@ describe('SettingsService (main)', () => { }); it('rejects an unsafe registered project before cleanup starts', async () => { - loadDatabaseFile.mockResolvedValue({ + stubDb({ projects: [{ path: '/tmp/dbt-studio-home' }], settings: {}, }); @@ -221,7 +258,7 @@ describe('SettingsService (main)', () => { }); it('preserves the macOS credential verification detail and requires restart', async () => { - loadDatabaseFile.mockResolvedValue({ projects: [], settings: {} }); + stubDb({ projects: [], settings: {} }); ( SecureStorageService.clearAllCredentials as jest.Mock ).mockRejectedValueOnce( @@ -239,7 +276,7 @@ describe('SettingsService (main)', () => { }); it('fails when a running task cannot be cancelled', async () => { - loadDatabaseFile.mockResolvedValue({ projects: [], settings: {} }); + stubDb({ projects: [], settings: {} }); (TaskManagerService.cancelAll as jest.Mock).mockReturnValueOnce(1); await expect( @@ -251,7 +288,7 @@ describe('SettingsService (main)', () => { it('times out a shutdown operation instead of waiting indefinitely', async () => { jest.useFakeTimers(); - loadDatabaseFile.mockResolvedValue({ projects: [], settings: {} }); + stubDb({ projects: [], settings: {} }); (FlowfileService.stop as jest.Mock).mockReturnValueOnce( new Promise(() => { // Simulate a shutdown operation that never settles. From c92ffaf45d33830d3750f1b7f77d9637ac86b40d Mon Sep 17 00:00:00 2001 From: ailegion Date: Wed, 9 Sep 2026 12:06:33 +0200 Subject: [PATCH 2/4] fixed: e2e tests to use new database store approach --- e2e/fixtures/electron.fixture.ts | 30 +++- e2e/tests/projects/project-lifecycle.spec.ts | 142 +++++-------------- 2 files changed, 65 insertions(+), 107 deletions(-) diff --git a/e2e/fixtures/electron.fixture.ts b/e2e/fixtures/electron.fixture.ts index e4bfd77b..00f4da10 100644 --- a/e2e/fixtures/electron.fixture.ts +++ b/e2e/fixtures/electron.fixture.ts @@ -21,6 +21,15 @@ export type ElectronFixtures = { userData: string; /** Whether to automatically skip the setup wizard by seeding database.json */ autoSkipSetup: boolean; + /** + * Project names to seed into database.json before the app launches. Each + * gets a minimal project directory + dbt_project.yml created automatically. + * Use this instead of writing to database.json after the app is already + * running — the app only ever re-reads the file on its own operations, so + * an external write made while it's live has no defined way to be picked + * up short of restarting it. + */ + extraProjects: string[]; /** The Electron application instance */ electronApp: ElectronApplication; /** The main browser window */ @@ -33,6 +42,7 @@ export type ElectronFixtures = { export const test = base.extend({ // Default to skipping setup for convenience in most tests autoSkipSetup: [true, { option: true }], + extraProjects: [[], { option: true }], // Create isolated userData directory for each test // biome-ignore lint/complexity/noEmptyPattern: Playwright requires object destructuring @@ -59,7 +69,7 @@ export const test = base.extend({ }, // Launch Electron app - electronApp: async ({ userData, autoSkipSetup }, use) => { + electronApp: async ({ userData, autoSkipSetup, extraProjects }, use) => { // Helper to seed database if skipping setup if (autoSkipSetup) { // Create projects directory @@ -81,6 +91,22 @@ export const test = base.extend({ const dbPath = path.join(userData, 'database.json'); + const seededProjects = extraProjects.map((name) => { + const projectPath = path.join(userData, 'projects', name); + fs.mkdirSync(projectPath, { recursive: true }); + fs.writeFileSync( + path.join(projectPath, 'dbt_project.yml'), + `name: ${name}\nversion: 1.0.0\nconfig-version: 2\n`, + ); + return { + id: `${name}-id`, + name, + path: projectPath, + createdAt: new Date().toISOString(), + isExtracted: false, + }; + }); + const settings = { schemaVersion: CURRENT_SCHEMA_VERSION, settings: { @@ -96,7 +122,7 @@ export const test = base.extend({ dbtSampleDirectory: path.join(userData, 'dbt_sample'), sampleRosettaMainConf: path.join(userData, 'main.conf'), }, - projects: [], + projects: seededProjects, connections: [], }; fs.writeFileSync(dbPath, JSON.stringify(settings, null, 2)); diff --git a/e2e/tests/projects/project-lifecycle.spec.ts b/e2e/tests/projects/project-lifecycle.spec.ts index e92fe4ad..06111244 100644 --- a/e2e/tests/projects/project-lifecycle.spec.ts +++ b/e2e/tests/projects/project-lifecycle.spec.ts @@ -1,6 +1,4 @@ import { Page, ElectronApplication } from '@playwright/test'; -import * as fs from 'fs'; -import * as path from 'path'; import { test, expect } from '../../fixtures/electron.fixture'; import { ProjectSelectionPage } from '../../page-objects/screens/ProjectSelection'; import { AppHelper } from '../../helpers/app.helper'; @@ -68,119 +66,53 @@ test.describe('Project Lifecycle', () => { await expect(sidebar).toBeVisible(); }); - test('should open existing project', async ({ electronApp, userData }) => { - // Seed project manually since tests run in isolation - const projectPath = path.join(userData, 'projects', 'Test_Project'); - const dbPath = path.join(userData, 'database.json'); - - // Create project directory and minimal dbt_project.yml - if (!fs.existsSync(projectPath)) { - fs.mkdirSync(projectPath, { recursive: true }); - fs.writeFileSync( - path.join(projectPath, 'dbt_project.yml'), - 'name: Test_Project\nversion: 1.0.0\nconfig-version: 2\n', - ); - } - - // Update database.json - const db = JSON.parse(fs.readFileSync(dbPath, 'utf8')); - // Avoid duplicate seeding - if (!db.projects.some((p: any) => p.name === 'Test_Project')) { - db.projects.push({ - id: 'test-project-id', - name: 'Test_Project', - path: projectPath, - createdAt: new Date().toISOString(), - isExtracted: false, - }); - fs.writeFileSync(dbPath, JSON.stringify(db, null, 2)); - } + test.describe('with an existing project seeded', () => { + // Seeded into database.json before the app launches (see + // electron.fixture.ts) rather than written to the file mid-test — the + // app only re-reads database.json on its own operations, so a write + // made to it while the app is already running has no defined way to be + // observed short of restarting the app. + test.use({ extraProjects: ['Test_Project'] }); - const stableWindow = await findStableWindow(electronApp); - // Reload window to ensure renderer picks up the DB changes - await stableWindow.reload(); - await stableWindow.waitForLoadState('domcontentloaded'); + test('should open existing project', async ({ electronApp }) => { + const stableWindow = await findStableWindow(electronApp); + const projectSelection = new ProjectSelectionPage(stableWindow); - // Wait for project selection screen again - await stableWindow.waitForSelector('[data-testid="project-selection"]', { - timeout: 10000, - }); + await projectSelection.selectProject('Test_Project'); - // Re-attach console listener after reload if needed (Playwright usually keeps it on Page, but handle might change?) - // Actually finding stableWindow again might return same page object. + // Verify project details screen is shown (or main app) + const sidebar = stableWindow.locator('[data-testid="sidebar"]'); + await expect(sidebar).toBeVisible(); + }); - const projectSelection = new ProjectSelectionPage(stableWindow); + test('should delete a project', async ({ electronApp }) => { + const stableWindow = await findStableWindow(electronApp); - // Verify and Select - await projectSelection.selectProject('Test_Project'); + // Get project card for verification + const projectCard = stableWindow.locator( + '[data-testid="project-card-Test_Project"]', + ); - // Verify project details screen is shown (or main app) - const sidebar = stableWindow.locator('[data-testid="sidebar"]'); - await expect(sidebar).toBeVisible(); - }); + // Click options button + const optionsBtn = stableWindow.locator( + '[data-testid="project-options-Test_Project"]', + ); + await optionsBtn.click(); - test('should delete a project', async ({ electronApp, userData }) => { - // Seed project manually since tests run in isolation - const projectPath = path.join(userData, 'projects', 'Test_Project'); - const dbPath = path.join(userData, 'database.json'); - - // Create project directory and minimal dbt_project.yml - if (!fs.existsSync(projectPath)) { - fs.mkdirSync(projectPath, { recursive: true }); - fs.writeFileSync( - path.join(projectPath, 'dbt_project.yml'), - 'name: Test_Project\nversion: 1.0.0\nconfig-version: 2\n', + // Click delete option + const deleteOption = stableWindow.locator( + '[data-testid="context-menu-delete"]', ); - } - - // Update database.json - const db = JSON.parse(fs.readFileSync(dbPath, 'utf8')); - // Avoid duplicate seeding - if (!db.projects.some((p: any) => p.name === 'Test_Project')) { - db.projects.push({ - id: 'test-project-id', - name: 'Test_Project', - path: projectPath, - createdAt: new Date().toISOString(), - isExtracted: false, - }); - fs.writeFileSync(dbPath, JSON.stringify(db, null, 2)); - } + await deleteOption.click(); - const stableWindow = await findStableWindow(electronApp); - // Reload window to ensure renderer picks up the DB changes - await stableWindow.reload(); - await stableWindow.waitForLoadState('domcontentloaded'); + // Confirm deletion + const confirmBtn = stableWindow.locator( + '[data-testid="confirm-delete-btn"]', + ); + await confirmBtn.click(); - // Wait for project selection screen again - await stableWindow.waitForSelector('[data-testid="project-selection"]', { - timeout: 10000, + // Verify project is removed + await expect(projectCard).not.toBeVisible(); }); - - // Get project card for verification - const projectCard = stableWindow.locator( - '[data-testid="project-card-Test_Project"]', - ); - - // Click options button - const optionsBtn = stableWindow.locator( - '[data-testid="project-options-Test_Project"]', - ); - await optionsBtn.click(); - - // Click delete option - const deleteOption = stableWindow.locator( - '[data-testid="context-menu-delete"]', - ); - await deleteOption.click(); - - // Confirm deletion - const confirmBtn = stableWindow.locator( - '[data-testid="confirm-delete-btn"]', - ); - await confirmBtn.click(); - - // Verify project is removed - await expect(projectCard).not.toBeVisible(); }); }); From 02a23bc4a92df0b383fc1374e5a9dc97ed4970ef Mon Sep 17 00:00:00 2001 From: ailegion Date: Wed, 9 Sep 2026 12:40:44 +0200 Subject: [PATCH 3/4] fix: close downgrade data-loss and credential-cache leak in database store - Handle database.json written by a newer app build (e.g. after a downgrade): back up the file and pass its contents through without running them through the reconstructive migration whitelist, which was silently dropping any field this older build didn't recognize. - Make DatabaseStore.getField/getSnapshot return a deep clone instead of a live reference into the in-memory cache. Callers that enrich a returned connection with a secure-storage credential for immediate use (e.g. ConnectorsService.extractSchemaFromConnection) were mutating the cache itself, so the next unrelated write could persist plaintext credentials to database.json. - Add a structuredClone polyfill to the Jest jsdom test environment (jsdom 20 here predates native support, added in v21). --- src/main/database/migrations.ts | 69 +++++++++++++----- src/main/database/store.ts | 32 ++++++++- tests/unit/__setup__/jest.setup.ts | 10 +++ tests/unit/main/database/migrations.test.ts | 42 +++++++++++ tests/unit/main/database/store.test.ts | 79 +++++++++++++++++++++ 5 files changed, 212 insertions(+), 20 deletions(-) diff --git a/src/main/database/migrations.ts b/src/main/database/migrations.ts index f0bd4547..f1e0b534 100644 --- a/src/main/database/migrations.ts +++ b/src/main/database/migrations.ts @@ -33,6 +33,39 @@ export function pendingMigrations(fromVersion: number): Migration[] { .sort((a, b) => a.version - b.version); } +function toShape(rawInput: unknown): RawShape { + return typeof rawInput === 'object' && rawInput !== null + ? (rawInput as RawShape) + : {}; +} + +// Fills in safe defaults for every field this build reads/writes. +// `preserveUnknownKeys` controls whether fields this build doesn't +// recognize are kept (spread in) or dropped — migrate() drops them +// because it's reconstructing a known-old shape into the current one; +// passthroughNewerVersion() keeps them because it has no business +// deleting fields a *newer* build understands and this one doesn't. +function withSafeDefaults( + shape: RawShape, + schemaVersion: number, + preserveUnknownKeys: boolean, +): DataBase { + return { + ...(preserveUnknownKeys ? shape : {}), + schemaVersion, + projects: (shape.projects as DataBase['projects']) ?? [], + settings: + (shape.settings as DataBase['settings']) ?? ({} as DataBase['settings']), + selectedProject: shape.selectedProject as DataBase['selectedProject'], + queries: (shape.queries as DataBase['queries']) ?? {}, + savedQueries: shape.savedQueries as DataBase['savedQueries'], + connections: (shape.connections as DataBase['connections']) ?? [], + sources: (shape.sources as DataBase['sources']) ?? [], + recentItems: (shape.recentItems as DataBase['recentItems']) ?? [], + icebergInstances: shape.icebergInstances as DataBase['icebergInstances'], + } as DataBase; +} + /** * Runs every pending migration in order and defensively fills in any field * still missing afterwards (covers installs older than the oldest @@ -41,27 +74,29 @@ export function pendingMigrations(fromVersion: number): Migration[] { * should reach the store's recovery path instead of this function. */ export function migrate(rawInput: unknown): DataBase { - let shape: RawShape = - typeof rawInput === 'object' && rawInput !== null - ? (rawInput as RawShape) - : {}; + let shape = toShape(rawInput); const startVersion = shape.schemaVersion ?? 0; pendingMigrations(startVersion).forEach((m) => { shape = m.migrate(shape); }); - return { - schemaVersion: CURRENT_SCHEMA_VERSION, - projects: (shape.projects as DataBase['projects']) ?? [], - settings: - (shape.settings as DataBase['settings']) ?? ({} as DataBase['settings']), - selectedProject: shape.selectedProject as DataBase['selectedProject'], - queries: (shape.queries as DataBase['queries']) ?? {}, - savedQueries: shape.savedQueries as DataBase['savedQueries'], - connections: (shape.connections as DataBase['connections']) ?? [], - sources: (shape.sources as DataBase['sources']) ?? [], - recentItems: (shape.recentItems as DataBase['recentItems']) ?? [], - icebergInstances: shape.icebergInstances as DataBase['icebergInstances'], - }; + return withSafeDefaults(shape, CURRENT_SCHEMA_VERSION, false); +} + +/** + * For a file written by a newer build than this one (e.g. after a + * downgrade): unlike migrate(), never reconstructs the object from a fixed + * key whitelist — any top-level field this build doesn't recognize + * survives untouched, and the recorded schemaVersion is left as whatever + * the newer build wrote. This build only acts on the fields it knows; + * anything it doesn't stays intact so an eventual re-upgrade loses nothing. + */ +export function passthroughNewerVersion(rawInput: unknown): DataBase { + const shape = toShape(rawInput); + return withSafeDefaults( + shape, + shape.schemaVersion ?? CURRENT_SCHEMA_VERSION, + true, + ); } diff --git a/src/main/database/store.ts b/src/main/database/store.ts index 3bfcb2e6..5708bd1d 100644 --- a/src/main/database/store.ts +++ b/src/main/database/store.ts @@ -1,7 +1,11 @@ import fs from 'fs'; import path from 'path'; import { DataBase } from '../../types/backend'; -import { CURRENT_SCHEMA_VERSION, migrate } from './migrations'; +import { + CURRENT_SCHEMA_VERSION, + migrate, + passthroughNewerVersion, +} from './migrations'; function defaultDatabase(): DataBase { return { @@ -43,10 +47,19 @@ export class DatabaseStore { this.cache = null; } + /** + * Returns a deep clone, never a live reference into the cache. Callers + * routinely enrich what they get back in place (e.g. materializing a + * secure-storage credential onto a connection object before using it) — + * with a cached store, mutating a live reference would poison the cache + * itself, and the next unrelated write would persist that mutation to + * disk. Cloning here is the one place that has to hold for every caller, + * present and future, rather than trusting each call site to remember. + */ async getField(key: K): Promise { return this.enqueue(async () => { const db = await this.ensureLoaded(); - return db[key]; + return structuredClone(db[key]); }); } @@ -55,9 +68,10 @@ export class DatabaseStore { * separate getField calls whenever a caller combines two or more fields * (e.g. joining projects to connections), since two separate getField * calls are two separate queue turns and a write could land in between. + * Also a deep clone, for the same reason as getField. */ async getSnapshot(): Promise> { - return this.enqueue(async () => this.ensureLoaded()); + return this.enqueue(async () => structuredClone(await this.ensureLoaded())); } async updateField( @@ -134,6 +148,18 @@ export class DatabaseStore { const onDiskVersion = (parsed as { schemaVersion?: number })?.schemaVersion ?? 0; + + if (onDiskVersion > CURRENT_SCHEMA_VERSION) { + // Written by a newer build than this one (e.g. after a downgrade). + // Back up first — this branch never persists on its own, but the + // very next write from this build otherwise would, with no backup + // ever having been taken — then read it without running it through + // migrate()'s reconstructive whitelist, which would silently drop + // every field this older build doesn't recognize. + await this.backup(raw, onDiskVersion); + return passthroughNewerVersion(parsed); + } + const needsMigration = onDiskVersion < CURRENT_SCHEMA_VERSION; if (needsMigration) { await this.backup(raw, onDiskVersion); diff --git a/tests/unit/__setup__/jest.setup.ts b/tests/unit/__setup__/jest.setup.ts index fca8b9d4..473deefe 100644 --- a/tests/unit/__setup__/jest.setup.ts +++ b/tests/unit/__setup__/jest.setup.ts @@ -19,5 +19,15 @@ if (!(global as any).TransformStream) { (global as any).TransformStream = TransformStream; } +// jsdom (v20 here) doesn't have structuredClone (added in jsdom v21). +// Production runs in the Electron main process (plain Node), where the +// real global is always present. +if (!(global as any).structuredClone) { + // eslint-disable-next-line global-require + const v8 = require('v8'); + (global as any).structuredClone = (value: unknown) => + v8.deserialize(v8.serialize(value)); +} + (global as any).fetch = jest.fn(); process.env.NODE_ENV = 'test'; diff --git a/tests/unit/main/database/migrations.test.ts b/tests/unit/main/database/migrations.test.ts index d515f696..e65baf7b 100644 --- a/tests/unit/main/database/migrations.test.ts +++ b/tests/unit/main/database/migrations.test.ts @@ -1,6 +1,7 @@ import { CURRENT_SCHEMA_VERSION, migrate, + passthroughNewerVersion, pendingMigrations, } from '../../../../src/main/database/migrations'; @@ -62,6 +63,47 @@ describe('migrations', () => { }); }); + describe('passthroughNewerVersion', () => { + it('preserves fields this build does not recognize instead of dropping them', () => { + const fromNewerBuild = { + schemaVersion: CURRENT_SCHEMA_VERSION + 1, + projects: [{ id: '1' }], + settings: {}, + queries: {}, + connections: [], + sources: [], + recentItems: [], + // A field only the newer build understands. + futureFeatureConfig: { enabled: true }, + }; + + const result = passthroughNewerVersion( + fromNewerBuild, + ) as unknown as typeof fromNewerBuild; + + expect(result.futureFeatureConfig).toEqual({ enabled: true }); + expect(result.projects).toEqual(fromNewerBuild.projects); + }); + + it('does not down-stamp schemaVersion to the current build version', () => { + const fromNewerBuild = { schemaVersion: CURRENT_SCHEMA_VERSION + 5 }; + + const result = passthroughNewerVersion(fromNewerBuild); + + expect(result.schemaVersion).toBe(CURRENT_SCHEMA_VERSION + 5); + }); + + it('still fills in safe defaults for known fields that are missing', () => { + const result = passthroughNewerVersion({ + schemaVersion: CURRENT_SCHEMA_VERSION + 1, + }); + + expect(result.projects).toEqual([]); + expect(result.connections).toEqual([]); + expect(result.recentItems).toEqual([]); + }); + }); + describe('pendingMigrations', () => { it('returns every migration above the given version, in ascending order', () => { const pending = pendingMigrations(0); diff --git a/tests/unit/main/database/store.test.ts b/tests/unit/main/database/store.test.ts index 809397be..5c4e0cec 100644 --- a/tests/unit/main/database/store.test.ts +++ b/tests/unit/main/database/store.test.ts @@ -94,6 +94,50 @@ describe('DatabaseStore', () => { legacyContent, ); }); + + it('backs up and preserves unknown fields from a file written by a newer build (downgrade)', async () => { + const newerVersion = CURRENT_SCHEMA_VERSION + 1; + const fromNewerBuild = JSON.stringify({ + schemaVersion: newerVersion, + projects: [{ id: 'p1' }], + settings: {}, + queries: {}, + connections: [], + sources: [], + recentItems: [], + // A field only the newer build understands — must survive. + futureFeatureConfig: { enabled: true }, + }); + fs.writeFileSync(dbPath, fromNewerBuild); + + const store = new DatabaseStore(dbPath); + const projects = await store.getField('projects'); + expect(projects).toEqual([{ id: 'p1' }]); + + // A backup was taken before this (older) build touched anything. + const backupFiles = fs + .readdirSync(dir) + .filter((name) => name.includes(`.v${newerVersion}.bak-`)); + expect(backupFiles).toHaveLength(1); + expect(fs.readFileSync(path.join(dir, backupFiles[0]), 'utf8')).toBe( + fromNewerBuild, + ); + + // The mere read did not touch the live file — same as any other + // read-only access. + const stillOnDisk = JSON.parse(fs.readFileSync(dbPath, 'utf8')); + expect(stillOnDisk.schemaVersion).toBe(newerVersion); + + // A subsequent write from this older build must not delete the + // field it doesn't understand. + await store.updateField('projects', (current) => [ + ...current, + { id: 'p2' } as never, + ]); + const afterWrite = JSON.parse(fs.readFileSync(dbPath, 'utf8')); + expect(afterWrite.futureFeatureConfig).toEqual({ enabled: true }); + expect(afterWrite.projects).toEqual([{ id: 'p1' }, { id: 'p2' }]); + }); }); describe('corruption recovery', () => { @@ -142,6 +186,41 @@ describe('DatabaseStore', () => { }); }); + describe('getField isolation', () => { + it('returns a value that is safe to mutate without affecting the store', async () => { + const store = new DatabaseStore(dbPath); + await store.updateField('connections', () => [ + { id: 'c1', connection: { type: 'postgres' } } as never, + ]); + + const connections = await store.getField('connections'); + // Simulates materializing a secure-storage credential onto a + // connection object for immediate use — a real, existing pattern in + // ConnectorsService (extractSchemaFromConnection, etc.). + (connections[0] as any).connection.password = 'super-secret-password'; + + const reread = await store.getField('connections'); + expect((reread[0] as any).connection.password).toBeUndefined(); + }); + + it('does not leak a mutated value into a later unrelated write', async () => { + const store = new DatabaseStore(dbPath); + await store.updateField('connections', () => [ + { id: 'c1', connection: { type: 'postgres' } } as never, + ]); + + const connections = await store.getField('connections'); + (connections[0] as any).connection.password = 'super-secret-password'; + + // An unrelated write (e.g. adding a recent item) must not carry the + // credential injected above into what actually gets persisted. + await store.updateField('recentItems', () => [{ id: 'r1' } as never]); + + const onDisk = JSON.parse(fs.readFileSync(dbPath, 'utf8')); + expect(onDisk.connections[0].connection.password).toBeUndefined(); + }); + }); + describe('updateField', () => { it('persists the new value and later reads observe it', async () => { const store = new DatabaseStore(dbPath); From c1f73f5920fdafd0433d749829233969c58a3d2b Mon Sep 17 00:00:00 2001 From: ailegion Date: Tue, 15 Sep 2026 12:56:16 +0200 Subject: [PATCH 4/4] fixed: backup import merges now computed inside updateField against current store value fixed: connection name uniqueness always enforced, removed allowReservedNames bypass --- src/main/services/backup.service.ts | 122 +++++++++++------------- src/main/services/connectors.service.ts | 20 ++-- 2 files changed, 62 insertions(+), 80 deletions(-) diff --git a/src/main/services/backup.service.ts b/src/main/services/backup.service.ts index 2766b320..59958d19 100644 --- a/src/main/services/backup.service.ts +++ b/src/main/services/backup.service.ts @@ -452,9 +452,6 @@ export default class BackupService { onProgress?.(doneSections, totalSections); }; - // ── Load current state ─────────────────────────────────────────────────── - const currentDb = await databaseStore.getSnapshot(); - // ── Helper to read JSON entry from ZIP (handles individual or legacy snapshot) ── const dbSnapshotEntry = zip.getEntry(DB_SNAPSHOT_ENTRY); let legacySnapshot: Record | null = null; @@ -497,17 +494,14 @@ export default class BackupService { 'connections.json', ); if (Array.isArray(connectionsData)) { - const existing = new Set( - (currentDb.connections || []).map((c: any) => c.id), - ); - const toAdd = connectionsData.filter((c: any) => !existing.has(c.id)); - if (toAdd.length > 0) { - await databaseStore.updateField('connections', (current) => [ - ...(current || []), - ...toAdd, - ]); - } - result.imported.connections = toAdd.length; + let importedCount = 0; + await databaseStore.updateField('connections', (current) => { + const existing = new Set((current || []).map((c: any) => c.id)); + const toAdd = connectionsData.filter((c: any) => !existing.has(c.id)); + importedCount = toAdd.length; + return [...(current || []), ...toAdd]; + }); + result.imported.connections = importedCount; } reportProgress(); } @@ -516,17 +510,14 @@ export default class BackupService { if (manifest.categories.includes('sources')) { const sourcesData = readCategoryJson('sources', 'sources.json'); if (Array.isArray(sourcesData)) { - const existing = new Set( - (currentDb.sources || []).map((s: any) => s.id), - ); - const toAdd = sourcesData.filter((s: any) => !existing.has(s.id)); - if (toAdd.length > 0) { - await databaseStore.updateField('sources', (current) => [ - ...(current || []), - ...toAdd, - ]); - } - result.imported.sources = toAdd.length; + let importedCount = 0; + await databaseStore.updateField('sources', (current) => { + const existing = new Set((current || []).map((s: any) => s.id)); + const toAdd = sourcesData.filter((s: any) => !existing.has(s.id)); + importedCount = toAdd.length; + return [...(current || []), ...toAdd]; + }); + result.imported.sources = importedCount; } reportProgress(); } @@ -538,26 +529,25 @@ export default class BackupService { 'savedQueries.json', ); if (savedQueriesData && typeof savedQueriesData === 'object') { - const mergedSavedQueries = { ...(currentDb.savedQueries || {}) }; let importedCount = 0; - - for (const [connId, queries] of Object.entries(savedQueriesData)) { - if (Array.isArray(queries)) { - const existing = mergedSavedQueries[connId] || []; - const existingIds = new Set(existing.map((q: any) => q.id)); - const toAdd = queries.filter((q: any) => !existingIds.has(q.id)); - if (toAdd.length > 0) { - mergedSavedQueries[connId] = [...existing, ...toAdd]; - importedCount += toAdd.length; + await databaseStore.updateField('savedQueries', (current) => { + const merged = { ...(current || {}) }; + importedCount = 0; + for (const [connId, queries] of Object.entries(savedQueriesData)) { + if (Array.isArray(queries)) { + const existing = merged[connId] || []; + const existingIds = new Set(existing.map((q: any) => q.id)); + const toAdd = queries.filter((q: any) => !existingIds.has(q.id)); + if (toAdd.length > 0) { + merged[connId] = [...existing, ...toAdd]; + importedCount += toAdd.length; + } } } - } + return merged; + }); if (importedCount > 0) { - await databaseStore.updateField('savedQueries', (current) => ({ - ...(current || {}), - ...mergedSavedQueries, - })); result.imported.savedQueries = importedCount; } } @@ -638,17 +628,14 @@ export default class BackupService { // 1. Iceberg const icebergData = readCategoryJson('datalake', 'icebergInstances.json'); if (Array.isArray(icebergData)) { - const existing = new Set( - (currentDb.icebergInstances || []).map((i: any) => i.id), - ); - const toAdd = icebergData.filter((i: any) => !existing.has(i.id)); - if (toAdd.length > 0) { - await databaseStore.updateField('icebergInstances', (current) => [ - ...(current || []), - ...toAdd, - ]); - } - importedCount += toAdd.length; + let icebergCount = 0; + await databaseStore.updateField('icebergInstances', (current) => { + const existing = new Set((current || []).map((i: any) => i.id)); + const toAdd = icebergData.filter((i: any) => !existing.has(i.id)); + icebergCount = toAdd.length; + return [...(current || []), ...toAdd]; + }); + importedCount += icebergCount; } // 2. DuckLake @@ -715,14 +702,12 @@ export default class BackupService { if (manifest.categories.includes('projects')) { const projectsData = readCategoryJson('projects', 'projects.json'); if (Array.isArray(projectsData)) { + const currentSettings = await databaseStore.getField('settings'); const targetBaseDir = - currentDb.settings?.projectsDirectory || + currentSettings?.projectsDirectory || path.join(os.homedir(), 'rosetta-dbt-studio-projects'); - const existingNames = new Set( - (currentDb.projects || []).map((p: any) => p.name), - ); - const newProjects: any[] = []; + const restoredProjects: any[] = []; for (const proj of projectsData) { // ZIP entry prefix always uses forward slashes (ZIP spec) @@ -794,18 +779,23 @@ export default class BackupService { // Always update the path in the restored metadata to the local destination proj.path = path.normalize(projectDest); - - if (!existingNames.has(proj.name)) { - newProjects.push(proj); - } + restoredProjects.push(proj); } - if (newProjects.length > 0) { - await databaseStore.updateField('projects', (current) => [ - ...(current || []), - ...newProjects, - ]); - result.imported.projects = newProjects.length; + let importedCount = 0; + await databaseStore.updateField('projects', (current) => { + const existingNames = new Set( + (current || []).map((p: any) => p.name), + ); + const toAdd = restoredProjects.filter( + (p: any) => !existingNames.has(p.name), + ); + importedCount = toAdd.length; + return [...(current || []), ...toAdd]; + }); + + if (importedCount > 0) { + result.imported.projects = importedCount; } } reportProgress(); diff --git a/src/main/services/connectors.service.ts b/src/main/services/connectors.service.ts index 97324e67..369a1dfd 100644 --- a/src/main/services/connectors.service.ts +++ b/src/main/services/connectors.service.ts @@ -341,7 +341,6 @@ export default class ConnectorsService { */ static async saveNewConnectionForTemplate( connection: ConnectionInput, - allowReservedNames = true, ): Promise { const connectionId = uuidV4(); const newConnection: ConnectionModel = { @@ -351,12 +350,9 @@ export default class ConnectorsService { await databaseStore.updateField('connections', (current) => { const connections = current ?? []; - // Validate connection name with optional allowReservedNames flag const nameValidation = this.validateConnectionName( connection.name, connections, - undefined, - allowReservedNames, ); if (!nameValidation.isValid) { throw new Error(nameValidation.message); @@ -1751,7 +1747,6 @@ export default class ConnectorsService { name: string, existingConnections: ConnectionModel[], excludeId?: string, - allowReservedNames?: boolean, ): { isValid: boolean; message?: string } { // Check for empty name if (!name.trim()) { @@ -1761,15 +1756,12 @@ export default class ConnectorsService { }; } - // Check for uniqueness (case-insensitive). Template connections pass - // allowReservedNames so the reserved "DBT Connection" name is exempt. - const duplicateExists = - !allowReservedNames && - existingConnections.some( - (conn) => - conn.connection.name.toLowerCase().trim() === - name.toLowerCase().trim() && conn.id !== excludeId, - ); + // Check for uniqueness (case-insensitive) + const duplicateExists = existingConnections.some( + (conn) => + conn.connection.name.toLowerCase().trim() === + name.toLowerCase().trim() && conn.id !== excludeId, + ); if (duplicateExists) { return {