const Database = require('better-sqlite3'); const { Pool } = require('pg'); const path = require('path'); const fs = require('fs'); const config = require('./config'); let db = null; let pgPool = null; let dbType = 'sqlite'; const dbOps = { async init() { if (config.USE_POSTGRES) { await this.initPostgres(); } else { this.initSqlite(); } }, async initPostgres() { dbType = 'postgres'; console.log(`Initializing PostgreSQL connection to ${config.POSTGRES_HOST}:${config.POSTGRES_PORT}/${config.POSTGRES_DB}`); // Wait for PostgreSQL to be ready await this.waitForPostgres(); // Create database if it doesn't exist await this.createDatabaseIfNotExists(); // Create connection pool pgPool = new Pool({ host: config.POSTGRES_HOST, port: config.POSTGRES_PORT, user: config.POSTGRES_USER, password: config.POSTGRES_PASSWORD, database: config.POSTGRES_DB, max: 20, idleTimeoutMillis: 30000, connectionTimeoutMillis: 10000, }); // Test connection try { const client = await pgPool.connect(); await client.query('SELECT 1'); client.release(); console.log('PostgreSQL connection established'); } catch (error) { console.error('PostgreSQL connection failed:', error.message); console.log('Falling back to SQLite'); this.initSqlite(); return; } // Create cache table const client = await pgPool.connect(); try { await client.query(` CREATE TABLE IF NOT EXISTS cache ( url TEXT PRIMARY KEY, html TEXT NOT NULL, created_at BIGINT NOT NULL, expires_at BIGINT NOT NULL ); CREATE INDEX IF NOT EXISTS idx_expires_at ON cache(expires_at); `); await client.query('COMMIT'); } catch (error) { console.error('Error creating PostgreSQL tables:', error); } finally { client.release(); } }, async waitForPostgres(maxRetries = 30, retryDelay = 2000) { console.log(`Waiting for PostgreSQL at ${config.POSTGRES_HOST}:${config.POSTGRES_PORT}...`); for (let attempt = 0; attempt < maxRetries; attempt++) { try { const tempPool = new Pool({ host: config.POSTGRES_HOST, port: config.POSTGRES_PORT, user: config.POSTGRES_USER, password: config.POSTGRES_PASSWORD, database: 'postgres', connectionTimeoutMillis: 5000, }); const client = await tempPool.connect(); await client.query('SELECT 1'); client.release(); await tempPool.end(); console.log(`PostgreSQL is ready after ${attempt + 1} attempts`); return true; } catch (error) { if (attempt < maxRetries - 1) { console.log(`PostgreSQL not ready (attempt ${attempt + 1}/${maxRetries}): ${error.message}`); await new Promise(resolve => setTimeout(resolve, retryDelay)); } else { console.error(`PostgreSQL connection failed after ${maxRetries} attempts: ${error.message}`); return false; } } } return false; }, async createDatabaseIfNotExists() { try { const tempPool = new Pool({ host: config.POSTGRES_HOST, port: config.POSTGRES_PORT, user: config.POSTGRES_USER, password: config.POSTGRES_PASSWORD, database: 'postgres', }); const client = await tempPool.connect(); // Check if database exists const result = await client.query( 'SELECT 1 FROM pg_database WHERE datname = $1', [config.POSTGRES_DB] ); if (result.rows.length === 0) { console.log(`Creating database '${config.POSTGRES_DB}'...`); // Note: CREATE DATABASE cannot be run in a transaction await client.query(`CREATE DATABASE ${config.POSTGRES_DB}`); console.log(`Database '${config.POSTGRES_DB}' created successfully`); } else { console.log(`Database '${config.POSTGRES_DB}' already exists`); } client.release(); await tempPool.end(); return true; } catch (error) { console.error(`Error creating database: ${error.message}`); return false; } }, initSqlite() { dbType = 'sqlite'; console.log(`Initializing SQLite database at ${config.DB_PATH}`); // Ensure database directory exists const dbDir = path.dirname(config.DB_PATH); if (!fs.existsSync(dbDir)) { fs.mkdirSync(dbDir, { recursive: true }); } // Initialize SQLite database db = new Database(config.DB_PATH); db.pragma('journal_mode = WAL'); // Create cache table if it doesn't exist db.exec(` CREATE TABLE IF NOT EXISTS cache ( url TEXT PRIMARY KEY, html TEXT NOT NULL, created_at INTEGER NOT NULL, expires_at INTEGER NOT NULL ); CREATE INDEX IF NOT EXISTS idx_expires_at ON cache(expires_at); `); }, async getCached(url, now) { if (dbType === 'postgres') { const client = await pgPool.connect(); try { const result = await client.query( 'SELECT html FROM cache WHERE url = $1 AND expires_at > $2', [url, now] ); return result.rows.length > 0 ? result.rows[0].html : null; } finally { client.release(); } } else { const getCached = db.prepare('SELECT html FROM cache WHERE url = ? AND expires_at > ?'); const result = getCached.get(url, now); return result ? result.html : null; } }, async getCachedSeo(url, now) { const seoKey = `seo:${url}`; if (dbType === 'postgres') { const client = await pgPool.connect(); try { const result = await client.query( 'SELECT html FROM cache WHERE url = $1 AND expires_at > $2', [seoKey, now] ); if (result.rows.length > 0) { return JSON.parse(result.rows[0].html); } return null; } finally { client.release(); } } else { const getCached = db.prepare('SELECT html FROM cache WHERE url = ? AND expires_at > ?'); const result = getCached.get(seoKey, now); return result ? JSON.parse(result.html) : null; } }, async setCache(url, html, now, expiresAt) { if (dbType === 'postgres') { const client = await pgPool.connect(); try { await client.query( 'INSERT INTO cache (url, html, created_at, expires_at) VALUES ($1, $2, $3, $4) ON CONFLICT (url) DO UPDATE SET html = $2, created_at = $3, expires_at = $4', [url, html, now, expiresAt] ); } finally { client.release(); } } else { const setCache = db.prepare('INSERT OR REPLACE INTO cache (url, html, created_at, expires_at) VALUES (?, ?, ?, ?)'); setCache.run(url, html, now, expiresAt); } }, async setCacheSeo(url, seoData, now, expiresAt) { const seoKey = `seo:${url}`; const html = JSON.stringify(seoData); if (dbType === 'postgres') { const client = await pgPool.connect(); try { await client.query( 'INSERT INTO cache (url, html, created_at, expires_at) VALUES ($1, $2, $3, $4) ON CONFLICT (url) DO UPDATE SET html = $2, created_at = $3, expires_at = $4', [seoKey, html, now, expiresAt] ); } finally { client.release(); } } else { const setCache = db.prepare('INSERT OR REPLACE INTO cache (url, html, created_at, expires_at) VALUES (?, ?, ?, ?)'); setCache.run(seoKey, html, now, expiresAt); } }, async getCachedMeta(url, now) { const metaKey = `meta:${url}`; if (dbType === 'postgres') { const client = await pgPool.connect(); try { const result = await client.query( 'SELECT html FROM cache WHERE url = $1 AND expires_at > $2', [metaKey, now] ); if (result.rows.length > 0) { return JSON.parse(result.rows[0].html); } return null; } finally { client.release(); } } else { const getCached = db.prepare('SELECT html FROM cache WHERE url = ? AND expires_at > ?'); const result = getCached.get(metaKey, now); return result ? JSON.parse(result.html) : null; } }, async setCacheMeta(url, metaData, now, expiresAt) { const metaKey = `meta:${url}`; const html = JSON.stringify(metaData); if (dbType === 'postgres') { const client = await pgPool.connect(); try { await client.query( 'INSERT INTO cache (url, html, created_at, expires_at) VALUES ($1, $2, $3, $4) ON CONFLICT (url) DO UPDATE SET html = $2, created_at = $3, expires_at = $4', [metaKey, html, now, expiresAt] ); } finally { client.release(); } } else { const setCache = db.prepare('INSERT OR REPLACE INTO cache (url, html, created_at, expires_at) VALUES (?, ?, ?, ?)'); setCache.run(metaKey, html, now, expiresAt); } }, async getCachedOutgoingCalls(url, now) { const callsKey = `outgoing-calls:${url}`; if (dbType === 'postgres') { const client = await pgPool.connect(); try { const result = await client.query( 'SELECT html FROM cache WHERE url = $1 AND expires_at > $2', [callsKey, now] ); if (result.rows.length > 0) { return JSON.parse(result.rows[0].html); } return null; } finally { client.release(); } } else { const getCached = db.prepare('SELECT html FROM cache WHERE url = ? AND expires_at > ?'); const result = getCached.get(callsKey, now); return result ? JSON.parse(result.html) : null; } }, async setCacheOutgoingCalls(url, callsData, now, expiresAt) { const callsKey = `outgoing-calls:${url}`; const html = JSON.stringify(callsData); if (dbType === 'postgres') { const client = await pgPool.connect(); try { await client.query( 'INSERT INTO cache (url, html, created_at, expires_at) VALUES ($1, $2, $3, $4) ON CONFLICT (url) DO UPDATE SET html = $2, created_at = $3, expires_at = $4', [callsKey, html, now, expiresAt] ); } finally { client.release(); } } else { const setCache = db.prepare('INSERT OR REPLACE INTO cache (url, html, created_at, expires_at) VALUES (?, ?, ?, ?)'); setCache.run(callsKey, html, now, expiresAt); } }, async getCachedResultingUrl(url, now) { const resultingUrlKey = `resulting-url:${url}`; if (dbType === 'postgres') { const client = await pgPool.connect(); try { const result = await client.query( 'SELECT html FROM cache WHERE url = $1 AND expires_at > $2', [resultingUrlKey, now] ); if (result.rows.length > 0) { return JSON.parse(result.rows[0].html); } return null; } finally { client.release(); } } else { const getCached = db.prepare('SELECT html FROM cache WHERE url = ? AND expires_at > ?'); const result = getCached.get(resultingUrlKey, now); return result ? JSON.parse(result.html) : null; } }, async setCacheResultingUrl(url, resultingUrlData, now, expiresAt) { const resultingUrlKey = `resulting-url:${url}`; const html = JSON.stringify(resultingUrlData); if (dbType === 'postgres') { const client = await pgPool.connect(); try { await client.query( 'INSERT INTO cache (url, html, created_at, expires_at) VALUES ($1, $2, $3, $4) ON CONFLICT (url) DO UPDATE SET html = $2, created_at = $3, expires_at = $4', [resultingUrlKey, html, now, expiresAt] ); } finally { client.release(); } } else { const setCache = db.prepare('INSERT OR REPLACE INTO cache (url, html, created_at, expires_at) VALUES (?, ?, ?, ?)'); setCache.run(resultingUrlKey, html, now, expiresAt); } }, async deleteExpired(now) { if (dbType === 'postgres') { const client = await pgPool.connect(); try { const result = await client.query( 'DELETE FROM cache WHERE expires_at <= $1', [now] ); return result.rowCount; } finally { client.release(); } } else { const deleteExpired = db.prepare('DELETE FROM cache WHERE expires_at <= ?'); const result = deleteExpired.run(now); return result.changes; } }, async deleteByUrl(url) { if (dbType === 'postgres') { const client = await pgPool.connect(); try { const result = await client.query( 'DELETE FROM cache WHERE url = $1', [url] ); return result.rowCount; } finally { client.release(); } } else { const deleteByUrl = db.prepare('DELETE FROM cache WHERE url = ?'); const result = deleteByUrl.run(url); return result.changes; } }, async deleteAll() { if (dbType === 'postgres') { const client = await pgPool.connect(); try { const result = await client.query('DELETE FROM cache'); return result.rowCount; } finally { client.release(); } } else { const result = db.prepare('DELETE FROM cache').run(); return result.changes; } }, async getStats() { const now = Date.now(); if (dbType === 'postgres') { const client = await pgPool.connect(); try { const totalResult = await client.query('SELECT COUNT(*) as total FROM cache'); const validResult = await client.query( 'SELECT COUNT(*) as valid FROM cache WHERE expires_at > $1', [now] ); return { total: parseInt(totalResult.rows[0].total, 10), valid: parseInt(validResult.rows[0].valid, 10), }; } finally { client.release(); } } else { const stats = db.prepare('SELECT COUNT(*) as total, COUNT(CASE WHEN expires_at > ? THEN 1 END) as valid FROM cache').get(now); return { total: stats.total, valid: stats.valid, }; } }, async close() { if (dbType === 'postgres' && pgPool) { await pgPool.end(); } else if (db) { db.close(); } }, getDbType() { return dbType; }, }; module.exports = dbOps;