/** * Database abstraction layer (async). * * Punkt 4: PostgreSQL for production, SQLite fallback for development. * Both modes expose the SAME async API: db.prepare(sql).run/get/all() return Promises. * * - If DATABASE_URL starts with 'postgresql://' → PostgreSQL (pg) * - Otherwise → SQLite (better-sqlite3, wrapped in Promises for unified async API) */ const path = require('path'); const DATABASE_URL = process.env.DATABASE_URL || ''; const usePostgres = DATABASE_URL.startsWith('postgresql://') || DATABASE_URL.startsWith('postgres://'); // Helper: convert SQLite ? placeholders to PostgreSQL $1, $2, etc. function convertPlaceholders(sql) { let idx = 0; return sql.replace(/\?/g, () => { idx++; return '$' + idx; }); } let db; if (usePostgres) { // ============ PostgreSQL mode (production) ============ const { Pool } = require('pg'); const pool = new Pool({ connectionString: DATABASE_URL, max: 10, idleTimeoutMillis: 30000, connectionTimeoutMillis: 10000, }); pool.on('error', (err) => { console.error('[DB] PostgreSQL Pool-Fehler:', err.message); }); console.log('[DB] PostgreSQL-Verbindung hergestellt (Production-Modus).'); db = { _pool: pool, _type: 'postgres', prepare(sql) { const pgSql = convertPlaceholders(sql); return { run: (...params) => { // For INSERT statements, append RETURNING id to get the generated ID const isInsert = pgSql.trim().toUpperCase().startsWith('INSERT'); const finalSql = isInsert && !pgSql.toUpperCase().includes('RETURNING') ? pgSql.replace(/;?\s*$/, ' RETURNING id') : pgSql; return pool.query(finalSql, params).then(result => ({ changes: result.rowCount, lastInsertRowid: result.rows[0]?.id || null, })); }, get: (...params) => pool.query(pgSql, params).then(result => result.rows[0] || null), all: (...params) => pool.query(pgSql, params).then(result => result.rows), }; }, exec(sql) { return pool.query(sql); }, pragma(_str) { return Promise.resolve({}); }, transaction(fn) { // H2: Transaktions-Contract — fn erhält eine txDb-Instanz, deren // prepare()-Statements auf DERSELBEN Connection laufen wie die Transaktion. // WICHTIG: Alle Statements innerhalb von fn MÜSSEN über txDb.prepare() // erzeugt werden. Statements, die außerhalb (über das globale db-Objekt) // vorbereitet wurden, laufen außerhalb der Transaktion und machen sie // wirkungslos (kein Rollback bei Fehlern). return async (...args) => { const client = await pool.connect(); try { await client.query('BEGIN'); const txDb = { _type: 'postgres', _inTransaction: true, prepare(sql) { const pgSql = convertPlaceholders(sql); return { run: (...params) => client.query(pgSql, params).then(result => ({ changes: result.rowCount, lastInsertRowid: result.rows[0]?.id || null, })), get: (...params) => client.query(pgSql, params).then(result => result.rows[0] || null), all: (...params) => client.query(pgSql, params).then(result => result.rows), }; }, exec: (sql) => client.query(sql), pragma: () => Promise.resolve({}), }; const result = await fn.call(txDb, ...args); await client.query('COMMIT'); return result; } catch (err) { await client.query('ROLLBACK'); throw err; } finally { client.release(); } }; }, close() { return pool.end(); }, }; } else { // ============ SQLite mode (development) ============ // Wrapped in Promises so the API is identical to PostgreSQL (async) const Database = require('better-sqlite3'); const dbPath = path.join(__dirname, 'data', 'workflow.db'); const sqliteDb = new Database(dbPath); sqliteDb.pragma('journal_mode = WAL'); sqliteDb.pragma('foreign_keys = ON'); console.log('[DB] SQLite-Datenbank verbunden (better-sqlite3, WAL-Modus, async-Wrapper).'); // H2: GLOBALER Transaktions-Mutex — SQLite hat nur EINE Connection, daher // dürfen sich nie zwei Transaktionen überlappen (auch nicht verschiedene // transaction()-Wrapper wie createTask + updateTemplate). Ein pro-Wrapper // Mutex würde "cannot start a transaction within a transaction" ermöglichen. let sqliteTxQueue = Promise.resolve(); db = { _type: 'sqlite', prepare(sql) { const stmt = sqliteDb.prepare(sql); return { run: (...params) => Promise.resolve(stmt.run(...params)), get: (...params) => Promise.resolve(stmt.get(...params)), all: (...params) => Promise.resolve(stmt.all(...params)), }; }, exec(sql) { sqliteDb.exec(sql); return Promise.resolve(); }, pragma(str) { sqliteDb.pragma(str); return Promise.resolve({}); }, transaction(fn) { // H2: Transaktions-Contract (siehe PostgreSQL-Modus) — fn erhält eine // txDb-Instanz, deren Statements innerhalb der Transaktion laufen. // better-sqlite3's transaction() ist synchron; die async fn wird über // den GLOBALEN Warteschlangen-Mutex serialisiert (siehe oben). return (...args) => { const run = async () => { const txDb = { _type: 'sqlite', _inTransaction: true, prepare(sql) { const stmt = sqliteDb.prepare(sql); return { run: (...params) => Promise.resolve(stmt.run(...params)), get: (...params) => Promise.resolve(stmt.get(...params)), all: (...params) => Promise.resolve(stmt.all(...params)), }; }, exec(sql) { sqliteDb.exec(sql); return Promise.resolve(); }, pragma(str) { sqliteDb.pragma(str); return Promise.resolve({}); }, }; // BEGIN IMMEDIATE sichert den Schreib-Lock für die gesamte Transaktion. // busy_timeout verhindert SQLITE_BUSY bei konkurrierenden Lesern (WAL). sqliteDb.pragma('busy_timeout = 5000'); sqliteDb.exec('BEGIN IMMEDIATE'); try { const result = await fn.call(txDb, ...args); sqliteDb.exec('COMMIT'); return result; } catch (err) { try { sqliteDb.exec('ROLLBACK'); } catch (rollbackErr) { console.error('[DB] Rollback-Fehler:', rollbackErr.message); } throw err; } }; const result = sqliteTxQueue.then(run, run); // Queue darf nie in einem Fehlerzustand hängen bleiben sqliteTxQueue = result.catch(() => {}); return result; }; }, close() { sqliteDb.close(); return Promise.resolve(); }, }; } module.exports = db;