database.js 5.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226
  1. const mysql = require('mysql2/promise');
  2. const sqlite3 = require('sqlite3').verbose();
  3. const { Sequelize } = require('sequelize');
  4. const { logger } = require('../utils/logger');
  5. let sequelize = null;
  6. let mysqlConnection = null;
  7. let sqliteDb = null;
  8. const connectDatabase = async () => {
  9. try {
  10. const dbType = process.env.DB_TYPE || 'postgres';
  11. switch (dbType) {
  12. case 'mysql':
  13. await connectMySQL();
  14. break;
  15. case 'postgres':
  16. await connectPostgreSQL();
  17. break;
  18. case 'sqlite':
  19. default:
  20. await connectSQLite();
  21. break;
  22. }
  23. logger.info(`Database connected successfully: ${dbType}`);
  24. return true;
  25. } catch (error) {
  26. logger.error('Database connection failed:', error);
  27. throw error;
  28. }
  29. };
  30. const connectMySQL = async () => {
  31. try {
  32. // Create connection pool
  33. mysqlConnection = mysql.createPool({
  34. host: process.env.DB_HOST || 'localhost',
  35. port: process.env.DB_PORT || 3306,
  36. user: process.env.DB_USER || 'root',
  37. password: process.env.DB_PASSWORD || '',
  38. database: process.env.DB_NAME || 'financial_data',
  39. waitForConnections: true,
  40. connectionLimit: 10,
  41. queueLimit: 0,
  42. acquireTimeout: 60000,
  43. timeout: 60000
  44. });
  45. // Test the connection
  46. const connection = await mysqlConnection.getConnection();
  47. await connection.ping();
  48. connection.release();
  49. logger.info('MySQL database connected successfully');
  50. } catch (error) {
  51. logger.error('MySQL connection failed:', error);
  52. throw error;
  53. }
  54. };
  55. // PostgreSQL connection function
  56. const connectPostgreSQL = async () => {
  57. try {
  58. sequelize = new Sequelize(
  59. process.env.DB_NAME || 'financial_data',
  60. process.env.DB_USER || 'postgres',
  61. process.env.DB_PASSWORD || 'mqldev@123',
  62. {
  63. host: process.env.DB_HOST || 'localhost',
  64. port: process.env.DB_PORT || 5432,
  65. dialect: 'postgres',
  66. pool: {
  67. max: 10,
  68. min: 0,
  69. acquire: 60000,
  70. idle: 10000
  71. },
  72. logging: process.env.NODE_ENV === 'development' ? logger.info.bind(logger) : false
  73. }
  74. );
  75. await sequelize.authenticate();
  76. logger.info('PostgreSQL database connected successfully');
  77. } catch (error) {
  78. logger.error('PostgreSQL connection failed:', error);
  79. throw error;
  80. }
  81. };
  82. const connectSQLite = async () => {
  83. return new Promise((resolve, reject) => {
  84. const dbFile = process.env.DB_FILE || './data/financial_data.db';
  85. sqliteDb = new sqlite3.Database(dbFile, (err) => {
  86. if (err) {
  87. logger.error('SQLite connection failed:', err);
  88. reject(err);
  89. } else {
  90. logger.info('SQLite database connected successfully');
  91. resolve();
  92. }
  93. });
  94. });
  95. };
  96. const getMySQLConnection = () => {
  97. if (!mysqlConnection) {
  98. throw new Error('MySQL connection not established');
  99. }
  100. return mysqlConnection;
  101. };
  102. const getSequelizeInstance = () => {
  103. if (!sequelize) {
  104. throw new Error('Sequelize instance not established');
  105. }
  106. return sequelize;
  107. };
  108. const getSQLiteInstance = () => {
  109. if (!sqliteDb) {
  110. throw new Error('SQLite database not established');
  111. }
  112. return sqliteDb;
  113. };
  114. const closeDatabase = async () => {
  115. try {
  116. if (mysqlConnection) {
  117. await mysqlConnection.end();
  118. logger.info('MySQL connection closed');
  119. }
  120. if (sequelize) {
  121. await sequelize.close();
  122. logger.info('PostgreSQL connection closed');
  123. }
  124. if (sqliteDb) {
  125. await new Promise((resolve) => {
  126. sqliteDb.close((err) => {
  127. if (err) {
  128. logger.error('Error closing SQLite database:', err);
  129. } else {
  130. logger.info('SQLite connection closed');
  131. }
  132. resolve();
  133. });
  134. });
  135. }
  136. } catch (error) {
  137. logger.error('Error closing database connections:', error);
  138. }
  139. };
  140. // Database query helpers
  141. const executeQuery = async (query, params = []) => {
  142. const dbType = process.env.DB_TYPE || 'postgres';
  143. switch (dbType) {
  144. case 'mysql':
  145. const [rows] = await mysqlConnection.execute(query, params);
  146. return rows;
  147. case 'postgres':
  148. const [results] = await sequelize.query(query, {
  149. replacements: params,
  150. type: Sequelize.QueryTypes.SELECT
  151. });
  152. return results;
  153. case 'sqlite':
  154. default:
  155. return new Promise((resolve, reject) => {
  156. sqliteDb.all(query, params, (err, rows) => {
  157. if (err) {
  158. reject(err);
  159. } else {
  160. resolve(rows);
  161. }
  162. });
  163. });
  164. }
  165. };
  166. const executeNonQuery = async (query, params = []) => {
  167. const dbType = process.env.DB_TYPE || 'postgres';
  168. switch (dbType) {
  169. case 'mysql':
  170. const [result] = await mysqlConnection.execute(query, params);
  171. return result;
  172. case 'postgres':
  173. const [updateResults] = await sequelize.query(query, {
  174. replacements: params,
  175. type: Sequelize.QueryTypes.UPDATE
  176. });
  177. return updateResults;
  178. case 'sqlite':
  179. default:
  180. return new Promise((resolve, reject) => {
  181. sqliteDb.run(query, params, function(err) {
  182. if (err) {
  183. reject(err);
  184. } else {
  185. resolve({ changes: this.changes, lastID: this.lastID });
  186. }
  187. });
  188. });
  189. }
  190. };
  191. module.exports = {
  192. connectDatabase,
  193. closeDatabase,
  194. executeQuery,
  195. executeNonQuery,
  196. getMySQLConnection,
  197. getSequelizeInstance,
  198. getSQLiteInstance
  199. };