From 16c3353b2197eae152bed7a21000f786caf9decf Mon Sep 17 00:00:00 2001 From: wayanrivan Date: Tue, 28 Jul 2026 11:57:25 +0700 Subject: [PATCH] update: add mutex in create bundle realtime data --- models/GpsTracksModels.js | 157 ++++++++++++++++++++++++++++++++++++-- package.json | 1 + 2 files changed, 150 insertions(+), 8 deletions(-) diff --git a/models/GpsTracksModels.js b/models/GpsTracksModels.js index 4ff30b4..6d6f95b 100755 --- a/models/GpsTracksModels.js +++ b/models/GpsTracksModels.js @@ -3,6 +3,16 @@ const MysqlHelpers = require(`../library/LibMysqlHelper`) const moment = require(`moment`) const fs = require('fs') const path = require('path') +var mutex = require('mutex') +var uuid = require('uuid') + +const rltmMutex = mutex({ + id: uuid.v4(), + strategy: { + name: 'redis', + connectionString: process.env.REDIS_URL || 'redis://127.0.0.1:6379', + }, +}) class GpsTracksModels { static DEFAULT_COUNTRY_ID = 1 @@ -134,13 +144,124 @@ class GpsTracksModels { }) } + // static bundleCreate2(logs = {}, rltm = {}) { + // // console.log(rltm.device_id + " : Start bundleCreate2") + // return new Promise(async (resolve, reject) => { + // try { + // const conn = await MysqlHelpers.createConnection() + // await MysqlHelpers.createTrx(conn) + // // console.log("createTrx : " + rltm.device_id) + + // 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], + // ) + // // console.log("insert t_gps_tracks : " + rltm.device_id) + + // 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`, + // ) + // } + // // console.log("insert tracks_${yy}${mm} : " + rltm.device_id) + // } + // } + + // if (rltmLength > 0 && typeof resLogs !== "undefined") + // rltm.master_id = resLogs.insertId + // if ( + // rltmLength > 0 && + // rltm.latitude !== null && + // rltm.longitude !== null + // ) { + // 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) + // // console.log("Commit bundleCreate2 : " + rltm.device_id) + // resolve({ + // type: "success", + // result: resLogs, + // }) + // } catch (err) { + // console.log("ERROR bundleCreate2 : " + rltm.device_id) + // reject(err) + // } + // }) + // } + static bundleCreate2(logs = {}, rltm = {}) { - // console.log(rltm.device_id + " : Start bundleCreate2") return new Promise(async (resolve, reject) => { + let lock = null try { const conn = await MysqlHelpers.createConnection() await MysqlHelpers.createTrx(conn) - // console.log("createTrx : " + rltm.device_id) let rltmLength = Object.keys(rltm).length let resLogs = undefined @@ -151,7 +272,6 @@ class GpsTracksModels { `INSERT INTO t_gps_tracks SET ?;`, [logs], ) - // console.log("insert t_gps_tracks : " + rltm.device_id) if (logs.action == "location") { const date = logs.crt_d @@ -163,9 +283,9 @@ class GpsTracksModels { await MysqlHelpers.queryTrx( conn, ` - INSERT INTO tracks_${yy}${mm} - SET ? - `, + INSERT INTO tracks_${yy}${mm} + SET ? + `, [logs], ) } catch (error) { @@ -175,17 +295,28 @@ class GpsTracksModels { JSON.stringify(logs) + `\n`, ) } - // console.log("insert tracks_${yy}${mm} : " + rltm.device_id) } } 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 = ?`, @@ -234,7 +365,7 @@ class GpsTracksModels { } await MysqlHelpers.commit(conn) - // console.log("Commit bundleCreate2 : " + rltm.device_id) + resolve({ type: "success", result: resLogs, @@ -242,6 +373,16 @@ class GpsTracksModels { } 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) + } + } } }) } diff --git a/package.json b/package.json index 520f3f6..56d2281 100755 --- a/package.json +++ b/package.json @@ -34,6 +34,7 @@ "jsonwebtoken": "^8.5.1", "moment": "^2.29.1", "morgan": "^1.10.0", + "mutex": "^1.0.4", "mysql": "^2.18.1", "mysql2": "^3.15.3", "nanoid": "^3.3.1",