/
Krams
/
nodemon2
Обзор
Документация
Войти
/
Krams
/
nodemon2
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
backend/database.js
316 строк
11 KB
Иван Новиков
final
25 май 2026, 17:46
25 май 2026, 17:46
2c76ffd
Код
Авторство
О чём код?
import sqlite3 from 'sqlite3'; import pg from 'pg'; import fs from 'fs'; sqlite3.verbose(); const { Pool } = pg; const SQLITE_DB_FILE = process.env.SQLITE_DB_FILE || './nodemon.db'; const DB_WRITE_MODE = String(process.env.DB_WRITE_MODE || 'sqlite').trim().toLowerCase(); const DB_READ_MODE = String(process.env.DB_READ_MODE || 'sqlite').trim().toLowerCase(); const DB_SHADOW_COMPARE = String(process.env.DB_SHADOW_COMPARE || '0').trim() === '1'; const DB_DUAL_WRITE_STRICT = String(process.env.DB_DUAL_WRITE_STRICT || '0').trim() === '1'; const DB_DISABLE_SQLITE_FALLBACK = String(process.env.DB_DISABLE_SQLITE_FALLBACK || '0').trim() === '1'; const sqliteFileExists = (() => { try { return fs.existsSync(SQLITE_DB_FILE); } catch { return false; } })(); const hasPgConfig = Boolean(process.env.PGHOST || process.env.PGDATABASE || process.env.DATABASE_URL); const pgPool = hasPgConfig ? new Pool({ connectionString: process.env.DATABASE_URL || undefined, host: process.env.PGHOST || undefined, port: process.env.PGPORT ? Number(process.env.PGPORT) : undefined, database: process.env.PGDATABASE || undefined, user: process.env.PGUSER || undefined, password: process.env.PGPASSWORD || undefined, ssl: process.env.PGSSLMODE === 'require' ? { rejectUnauthorized: false } : undefined, max: Number(process.env.PGPOOL_MAX || 10), }) : null; const sqliteIsRequired = DB_READ_MODE !== 'postgres' || DB_WRITE_MODE !== 'postgres' || (!DB_DISABLE_SQLITE_FALLBACK && DB_READ_MODE === 'postgres') || (DB_SHADOW_COMPARE && DB_READ_MODE === 'postgres'); const sqliteDb = sqliteIsRequired && sqliteFileExists ? new sqlite3.Database(SQLITE_DB_FILE) : null; if (sqliteIsRequired && !sqliteDb) { console.warn( `[DB] SQLite is required by config but file is missing: ${SQLITE_DB_FILE}. ` + `Either restore the file, set DB_SHADOW_COMPARE=0, or switch read/write modes to postgres-only with no fallback.` ); } console.log( [ '[DB] config:', `read=${DB_READ_MODE}`, `write=${DB_WRITE_MODE}`, `shadowCompare=${DB_SHADOW_COMPARE ? '1' : '0'}`, `dualWriteStrict=${DB_DUAL_WRITE_STRICT ? '1' : '0'}`, `disableSqliteFallback=${DB_DISABLE_SQLITE_FALLBACK ? '1' : '0'}`, `pgConfigured=${pgPool ? '1' : '0'}`, `sqliteConfigured=${sqliteDb ? '1' : '0'}`, `sqliteFile=${SQLITE_DB_FILE}`, ].join(' ') ); function transformSqlForPostgres(sql = '') { let normalized = String(sql); normalized = normalized.replace(/\?/g, () => `$${placeholderIdx++}`); normalized = normalized.replace(/INTEGER\s+PRIMARY\s+KEY\s+AUTOINCREMENT/gi, 'BIGSERIAL PRIMARY KEY'); normalized = normalized.replace(/datetime\('now',\s*'(-?\d+)\s+hours'\)/gi, (_, hours) => `NOW() + INTERVAL '${hours} hours'`); normalized = normalized.replace(/datetime\('now'\)/gi, 'NOW()'); normalized = normalized.replace(/INSERT\s+OR\s+IGNORE/gi, 'INSERT'); return normalized; } function isTransactionCommand(sql = '') { const safe = String(sql).trim().toUpperCase(); return safe === 'BEGIN' || safe === 'BEGIN TRANSACTION' || safe === 'COMMIT' || safe === 'ROLLBACK'; } function compareRows(a, b) { try { return JSON.stringify(a) === JSON.stringify(b); } catch { return false; } } let placeholderIdx = 0; function toPgSql(sql) { placeholderIdx = 1; return transformSqlForPostgres(sql); } function runSqlite(sql, params = []) { if (!sqliteDb) return Promise.reject(new Error('SQLite database is not configured')); return new Promise((resolve, reject) => { sqliteDb.run(sql, params, function onRun(err) { if (err) return reject(err); resolve({ lastID: this.lastID, changes: this.changes }); }); }); } async function runPg(sql, params = []) { if (!pgPool) throw new Error('PostgreSQL pool is not configured'); const baseSql = toPgSql(sql); const isInsert = /^\s*INSERT\s+/i.test(baseSql); const hasReturning = /\bRETURNING\b/i.test(baseSql); let result; if (isInsert && !hasReturning) { try { result = await pgPool.query(`${baseSql} RETURNING id`, params); } catch { result = await pgPool.query(baseSql, params); } } else { result = await pgPool.query(baseSql, params); } const firstRow = Array.isArray(result.rows) ? result.rows[0] : null; const maybeId = firstRow && Object.prototype.hasOwnProperty.call(firstRow, 'id') ? firstRow.id : undefined; return { lastID: maybeId, changes: result.rowCount || 0 }; } function getSqlite(sql, params = []) { if (!sqliteDb) return Promise.reject(new Error('SQLite database is not configured')); return new Promise((resolve, reject) => { sqliteDb.get(sql, params, (err, row) => { if (err) return reject(err); resolve(row || null); }); }); } async function getPg(sql, params = []) { if (!pgPool) throw new Error('PostgreSQL pool is not configured'); const result = await pgPool.query(toPgSql(sql), params); return result.rows?.[0] || null; } function allSqlite(sql, params = []) { if (!sqliteDb) return Promise.reject(new Error('SQLite database is not configured')); return new Promise((resolve, reject) => { sqliteDb.all(sql, params, (err, rows) => { if (err) return reject(err); resolve(rows || []); }); }); } async function allPg(sql, params = []) { if (!pgPool) throw new Error('PostgreSQL pool is not configured'); const result = await pgPool.query(toPgSql(sql), params); return result.rows || []; } function invokeRunCallback(callback, err, meta) { if (typeof callback !== 'function') return; const context = { lastID: meta?.lastID, changes: meta?.changes || 0, }; callback.call(context, err || null); } function invokeCallback(callback, err, payload) { if (typeof callback !== 'function') return; callback(err || null, payload); } const db = { serialize(fn) { if (typeof fn === 'function') fn(); }, run(sql, params, callback) { const args = Array.isArray(params) ? params : []; const cb = typeof params === 'function' ? params : callback; const shouldMirror = DB_WRITE_MODE === 'both'; const writePrimary = DB_WRITE_MODE === 'postgres' ? 'postgres' : 'sqlite'; const writeSecondary = writePrimary === 'sqlite' ? 'postgres' : 'sqlite'; const primaryRunner = writePrimary === 'postgres' ? runPg : runSqlite; const secondaryRunner = writeSecondary === 'postgres' ? runPg : runSqlite; const transactionCommand = isTransactionCommand(sql); primaryRunner(sql, args) .then(async (primaryMeta) => { if (shouldMirror && !transactionCommand) { const mirror = secondaryRunner(sql, args).catch((mirrorError) => { console.error('[DB] dual-write secondary failed:', mirrorError?.message || mirrorError); if (DB_DUAL_WRITE_STRICT) throw mirrorError; }); if (DB_DUAL_WRITE_STRICT) { await mirror; } } invokeRunCallback(cb, null, primaryMeta); }) .catch((error) => { invokeRunCallback(cb, error, null); }); return this; }, get(sql, params, callback) { const args = Array.isArray(params) ? params : []; const cb = typeof params === 'function' ? params : callback; const primaryRead = DB_READ_MODE === 'postgres' ? 'postgres' : 'sqlite'; const primaryReader = primaryRead === 'postgres' ? getPg : getSqlite; const shadowReader = primaryRead === 'postgres' ? getSqlite : getPg; primaryReader(sql, args) .then(async (row) => { if (DB_SHADOW_COMPARE && pgPool && !/PRAGMA/i.test(String(sql))) { try { const shadowRow = await shadowReader(sql, args); if (!compareRows(row, shadowRow)) { console.warn('[DB] shadow mismatch (get):', String(sql).slice(0, 140)); } } catch (shadowError) { console.warn('[DB] shadow compare failed (get):', shadowError?.message || shadowError); } } invokeCallback(cb, null, row); }) .catch(async (error) => { if (primaryRead === 'postgres') { if (DB_DISABLE_SQLITE_FALLBACK) { invokeCallback(cb, error, null); return; } try { const fallback = await getSqlite(sql, args); invokeCallback(cb, null, fallback); return; } catch { } } invokeCallback(cb, error, null); }); }, all(sql, params, callback) { const args = Array.isArray(params) ? params : []; const cb = typeof params === 'function' ? params : callback; const primaryRead = DB_READ_MODE === 'postgres' ? 'postgres' : 'sqlite'; const primaryReader = primaryRead === 'postgres' ? allPg : allSqlite; const shadowReader = primaryRead === 'postgres' ? allSqlite : allPg; primaryReader(sql, args) .then(async (rows) => { if (DB_SHADOW_COMPARE && pgPool && !/PRAGMA/i.test(String(sql))) { try { const shadowRows = await shadowReader(sql, args); if (!compareRows(rows, shadowRows)) { console.warn('[DB] shadow mismatch (all):', String(sql).slice(0, 140)); } } catch (shadowError) { console.warn('[DB] shadow compare failed (all):', shadowError?.message || shadowError); } } invokeCallback(cb, null, rows); }) .catch(async (error) => { if (primaryRead === 'postgres') { if (DB_DISABLE_SQLITE_FALLBACK) { invokeCallback(cb, error, []); return; } try { const fallback = await allSqlite(sql, args); invokeCallback(cb, null, fallback); return; } catch { } } invokeCallback(cb, error, []); }); }, async withTransaction(work) { if (typeof work !== 'function') { throw new Error('withTransaction expects a function'); } await runSqlite('BEGIN TRANSACTION'); try { const result = await work(db); await runSqlite('COMMIT'); return result; } catch (error) { await runSqlite('ROLLBACK'); throw error; } }, async hasColumn(tableName, columnName) { const safeTable = String(tableName || '').trim(); const safeColumn = String(columnName || '').trim(); if (!safeTable || !safeColumn) return false; if (DB_READ_MODE === 'postgres' && pgPool) { const row = await getPg( `SELECT column_name FROM information_schema.columns WHERE table_schema = 'public' AND table_name = ? AND column_name = ? LIMIT 1`, [safeTable, safeColumn] ); return Boolean(row); } const rows = await allSqlite(`PRAGMA table_info(${safeTable})`); return Array.isArray(rows) && rows.some((row) => row?.name === safeColumn); }, }; export default db;