| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226 |
- 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
- };
|