383 lines
13 KiB
JavaScript
383 lines
13 KiB
JavaScript
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 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;
|