const mysql = require('mysql2/promise'); const sqlite3 = require('sqlite3').verbose(); const { Sequelize } = require('sequelize'); const { logger } = require('../utils/logger'); let sequelize = null; let mysqlConnection = null; let sqliteDb = null; const connectDatabase = async () => { try { const dbType = process.env.DB_TYPE || 'postgres'; switch (dbType) { case 'mysql': await connectMySQL(); break; case 'postgres': await connectPostgreSQL(); break; case 'sqlite': default: await connectSQLite(); break; } logger.info(`Database connected successfully: ${dbType}`); return true; } catch (error) { logger.error('Database connection failed:', error); throw error; } }; const connectMySQL = async () => { try { // Create connection pool mysqlConnection = mysql.createPool({ host: process.env.DB_HOST || 'localhost', port: process.env.DB_PORT || 3306, user: process.env.DB_USER || 'root', password: process.env.DB_PASSWORD || '', database: process.env.DB_NAME || 'financial_data', waitForConnections: true, connectionLimit: 10, queueLimit: 0, acquireTimeout: 60000, timeout: 60000 }); // Test the connection const connection = await mysqlConnection.getConnection(); await connection.ping(); connection.release(); logger.info('MySQL database connected successfully'); } catch (error) { logger.error('MySQL connection failed:', error); throw error; } }; // PostgreSQL connection function const connectPostgreSQL = async () => { try { sequelize = new Sequelize( process.env.DB_NAME || 'financial_data', process.env.DB_USER || 'postgres', process.env.DB_PASSWORD || 'mqldev@123', { host: process.env.DB_HOST || 'localhost', port: process.env.DB_PORT || 5432, dialect: 'postgres', pool: { max: 10, min: 0, acquire: 60000, idle: 10000 }, logging: process.env.NODE_ENV === 'development' ? logger.info.bind(logger) : false } ); await sequelize.authenticate(); logger.info('PostgreSQL database connected successfully'); } catch (error) { logger.error('PostgreSQL connection failed:', error); throw error; } }; const connectSQLite = async () => { return new Promise((resolve, reject) => { const dbFile = process.env.DB_FILE || './data/financial_data.db'; sqliteDb = new sqlite3.Database(dbFile, (err) => { if (err) { logger.error('SQLite connection failed:', err); reject(err); } else { logger.info('SQLite database connected successfully'); resolve(); } }); }); }; const getMySQLConnection = () => { if (!mysqlConnection) { throw new Error('MySQL connection not established'); } return mysqlConnection; }; const getSequelizeInstance = () => { if (!sequelize) { throw new Error('Sequelize instance not established'); } return sequelize; }; const getSQLiteInstance = () => { if (!sqliteDb) { throw new Error('SQLite database not established'); } return sqliteDb; }; const closeDatabase = async () => { try { if (mysqlConnection) { await mysqlConnection.end(); logger.info('MySQL connection closed'); } if (sequelize) { await sequelize.close(); logger.info('PostgreSQL connection closed'); } if (sqliteDb) { await new Promise((resolve) => { sqliteDb.close((err) => { if (err) { logger.error('Error closing SQLite database:', err); } else { logger.info('SQLite connection closed'); } resolve(); }); }); } } catch (error) { logger.error('Error closing database connections:', error); } }; // Database query helpers const executeQuery = async (query, params = []) => { const dbType = process.env.DB_TYPE || 'postgres'; switch (dbType) { case 'mysql': const [rows] = await mysqlConnection.execute(query, params); return rows; case 'postgres': const [results] = await sequelize.query(query, { replacements: params, type: Sequelize.QueryTypes.SELECT }); return results; case 'sqlite': default: return new Promise((resolve, reject) => { sqliteDb.all(query, params, (err, rows) => { if (err) { reject(err); } else { resolve(rows); } }); }); } }; const executeNonQuery = async (query, params = []) => { const dbType = process.env.DB_TYPE || 'postgres'; switch (dbType) { case 'mysql': const [result] = await mysqlConnection.execute(query, params); return result; case 'postgres': const [updateResults] = await sequelize.query(query, { replacements: params, type: Sequelize.QueryTypes.UPDATE }); return updateResults; case 'sqlite': default: return new Promise((resolve, reject) => { sqliteDb.run(query, params, function(err) { if (err) { reject(err); } else { resolve({ changes: this.changes, lastID: this.lastID }); } }); }); } }; module.exports = { connectDatabase, closeDatabase, executeQuery, executeNonQuery, getMySQLConnection, getSequelizeInstance, getSQLiteInstance };