This commit is contained in:
@@ -0,0 +1,298 @@
|
||||
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 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 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;
|
||||
Reference in New Issue
Block a user