Files
workflow-app/backend/db.js

223 lines
7.7 KiB
JavaScript

/**
* 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) => {
// BUGFIX (v8): lastInsertRowid innerhalb der Transaktion —
// identische RETURNING-Logik wie im outer prepare(). Ohne
// RETURNING id liefert pg keine rows → lastInsertRowid war
// null → Folge-Inserts (task_values mit task_id=null) schlugen
// mit FK-Verstoß fehl (500 beim Abschicken einer Vorlage).
const isInsert = pgSql.trim().toUpperCase().startsWith('INSERT');
const finalSql = isInsert && !pgSql.toUpperCase().includes('RETURNING')
? pgSql.replace(/;?\s*$/, ' RETURNING id')
: pgSql;
return client.query(finalSql, 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;