diff --git a/library/LibMysqlHelper.js b/library/LibMysqlHelper.js index a7e0b1a..9869f37 100755 --- a/library/LibMysqlHelper.js +++ b/library/LibMysqlHelper.js @@ -2,7 +2,7 @@ const db = require(`../config/dbMysqlConn`); // const Promise = require("bluebird"); class MysqlHelpers { - static async insert (table, data) { + static async insert(table, data) { return new Promise((resolve, reject) => { const query = `INSERT INTO ${table} SET ?;`; db.getConnection(function (err, conn) { @@ -42,7 +42,7 @@ class MysqlHelpers { }); }) } - static async update (table, data, colId, valId) { + static async update(table, data, colId, valId) { return new Promise((resolve, reject) => { const query = `UPDATE ${table} SET ? WHERE ${colId} = ?;`; db.getConnection(function (err, conn) { @@ -81,7 +81,7 @@ class MysqlHelpers { }); }) } - static async delete (table, colId, valId) { + static async delete(table, colId, valId) { return new Promise((resolve, reject) => { const query = `DELETE FROM ${table} WHERE ${colId} = ?;`; db.getConnection(function (err, conn) { @@ -120,7 +120,7 @@ class MysqlHelpers { }); }) } - static async createConnection () { + static async createConnection() { return new Promise((resolve, reject) => { db.getConnection(function (err, conn) { if (err) { @@ -132,18 +132,18 @@ class MysqlHelpers { }); }) } - static async getDbMysqlConn () { + static async getDbMysqlConn() { return new Promise((resolve, reject) => { resolve(db); }) } - static async releaseConnection (conn) { + static async releaseConnection(conn) { return new Promise((resolve, reject) => { if (conn) conn.release(); resolve(true); }) } - static async createTrx (conn) { + static async createTrx(conn) { return new Promise((resolve, reject) => { conn.beginTransaction(async function (err) { if (err) { @@ -155,7 +155,7 @@ class MysqlHelpers { }); }) } - static async queryTrx (conn, query = '', params = []) { + static async queryTrx(conn, query = '', params = []) { return new Promise((resolve, reject) => { conn.query(query, params, function (err, result) { if (err) { @@ -170,7 +170,7 @@ class MysqlHelpers { }); }) } - static async query (conn, query = '', params = []) { + static async query(conn, query = '', params = []) { return new Promise((resolve, reject) => { conn.query(query, params, function (err, result) { if (err) return reject(err); @@ -179,7 +179,7 @@ class MysqlHelpers { }); }) } - static async commit (conn) { + static async commit(conn) { return new Promise((resolve, reject) => { conn.commit(async function (err) { if (err) { @@ -192,13 +192,24 @@ class MysqlHelpers { }); }) } - static async rollback (conn) { - return new Promise((resolve, reject) => { - conn.rollback(async function () { - conn.release(); - reject(err); - }); - return false; + // static async rollback (conn) { + // return new Promise((resolve, reject) => { + // conn.rollback(async function () { + // conn.release(); + // reject(err); + // }); + // return false; + // }) + // } + static async rollback(conn) { + return new Promise((resolve) => { + conn.rollback(() => { + try { + conn.release() + } catch (e) { } + + resolve(true) + }) }) } } diff --git a/models/GpsTracksModels.js b/models/GpsTracksModels.js index 6d6f95b..57b1a9d 100755 --- a/models/GpsTracksModels.js +++ b/models/GpsTracksModels.js @@ -256,6 +256,141 @@ class GpsTracksModels { // }) // } + /** MUTEX ESPECIALLY FOR RLTM */ + + // static bundleCreate2(logs = {}, rltm = {}) { + // return new Promise(async (resolve, reject) => { + // let lock = null + // try { + // const conn = await MysqlHelpers.createConnection() + // await MysqlHelpers.createTrx(conn) + + // let rltmLength = Object.keys(rltm).length + // let resLogs = undefined + + // if (Object.keys(logs).length > 0) { + // resLogs = await MysqlHelpers.queryTrx( + // conn, + // `INSERT INTO t_gps_tracks SET ?;`, + // [logs], + // ) + + // if (logs.action == "location") { + // const date = logs.crt_d + // const mm = moment.unix(date).format("MM") + // const yy = moment.unix(date).format("YY") + // logs.id = resLogs.insertId + + // try { + // await MysqlHelpers.queryTrx( + // conn, + // ` + // INSERT INTO tracks_${yy}${mm} + // SET ? + // `, + // [logs], + // ) + // } catch (error) { + // console.log(error) + // fs.appendFileSync( + // path.join(__dirname, `error_data_log.txt`), + // JSON.stringify(logs) + `\n`, + // ) + // } + // } + // } + + // if (rltmLength > 0 && typeof resLogs !== "undefined") + // rltm.master_id = resLogs.insertId + + // if ( + // rltmLength > 0 && + // rltm.latitude !== null && + // rltm.longitude !== null + // ) { + // // === KUNCI DI SINI === + // // Lock per device_id supaya SELECT-DELETE-INSERT untuk + // // device yang sama tidak berjalan bersamaan dari request/worker + // // lain. Inilah yang tadinya memicu ER_DUP_ENTRY dan + // // "Lock wait timeout" karena banyak transaksi menumpuk + // // menunggu baris yang sama. + // lock = await rltmMutex.lock(`t_gps_tracks_rltm:${rltm.device_id}`, { + // duration: 50000, // maksimal lock dipegang 5 detik + // maxWait: 50000, // menunggu lock maksimal 10 detik lalu error + // }) + + // let devices = await MysqlHelpers.queryTrx( + // conn, + // `SELECT id FROM t_gps_tracks_rltm as rltm WHERE rltm.device_id = ?`, + // [rltm.device_id], + // ) + // if (devices.length > 1) + // await MysqlHelpers.queryTrx( + // conn, + // `DELETE from t_gps_tracks_rltm WHERE device_id = ?;`, + // [rltm.device_id], + // ) + + // if (rltm.vhc_id != 0) { + // let vhcs = await MysqlHelpers.queryTrx( + // conn, + // `SELECT id FROM t_gps_tracks_rltm as rltm WHERE rltm.vhc_id = ?`, + // [rltm.vhc_id], + // ) + // if (vhcs.length > 1) + // await MysqlHelpers.queryTrx( + // conn, + // `DELETE from t_gps_tracks_rltm WHERE vhc_id = ?;`, + // [rltm.vhc_id], + // ) + // } + + // if (rltm.drv_id != 0) { + // let drvs = await MysqlHelpers.queryTrx( + // conn, + // `SELECT id FROM t_gps_tracks_rltm as rltm WHERE rltm.drv_id = ?`, + // [rltm.drv_id], + // ) + // if (drvs.length > 1) + // await MysqlHelpers.queryTrx( + // conn, + // `DELETE from t_gps_tracks_rltm WHERE drv_id = ?;`, + // [rltm.drv_id], + // ) + // } + + // await MysqlHelpers.queryTrx( + // conn, + // `INSERT INTO t_gps_tracks_rltm SET ? ON DUPLICATE KEY UPDATE ?;`, + // [rltm, rltm], + // ) + // } + + // await MysqlHelpers.commit(conn) + + // resolve({ + // type: "success", + // result: resLogs, + // }) + // } catch (err) { + // console.log("ERROR bundleCreate2 : " + rltm.device_id) + // reject(err) + // } finally { + // // Lock selalu dilepas, baik sukses maupun gagal, + // // supaya tidak deadlock request device yang sama berikutnya. + // if (lock) { + // try { + // await rltmMutex.unlock(lock) + // } catch (unlockErr) { + // console.log("ERROR unlock mutex : " + rltm.device_id, unlockErr) + // } + // } + // } + // }) + // } + + /** MUTEX FOR ALL REQUEST */ + static bundleCreate2(logs = {}, rltm = {}) { return new Promise(async (resolve, reject) => { let lock = null @@ -273,28 +408,38 @@ class GpsTracksModels { [logs], ) - if (logs.action == "location") { - const date = logs.crt_d - const mm = moment.unix(date).format("MM") - const yy = moment.unix(date).format("YY") - logs.id = resLogs.insertId + let trackLock = null + + try { + if (logs.action == "location") { + + trackLock = await rltmMutex.lock( + `tracks_insert:${logs.device_id}`, + { + duration: 90000, + maxWait: 90000, + } + ) + + const date = logs.crt_d + const mm = moment.unix(date).format("MM") + const yy = moment.unix(date).format("YY") + + logs.id = resLogs.insertId - try { await MysqlHelpers.queryTrx( conn, ` - INSERT INTO tracks_${yy}${mm} - SET ? - `, - [logs], - ) - } catch (error) { - console.log(error) - fs.appendFileSync( - path.join(__dirname, `error_data_log.txt`), - JSON.stringify(logs) + `\n`, + INSERT INTO tracks_${yy}${mm} + SET ? + `, + [logs] ) } + } finally { + if (trackLock) { + await rltmMutex.unlock(trackLock) + } } }