update: add mutex in create bundle realtime data
This commit is contained in:
@ -3,6 +3,16 @@ const MysqlHelpers = require(`../library/LibMysqlHelper`)
|
|||||||
const moment = require(`moment`)
|
const moment = require(`moment`)
|
||||||
const fs = require('fs')
|
const fs = require('fs')
|
||||||
const path = require('path')
|
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 {
|
class GpsTracksModels {
|
||||||
static DEFAULT_COUNTRY_ID = 1
|
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 = {}) {
|
static bundleCreate2(logs = {}, rltm = {}) {
|
||||||
// console.log(rltm.device_id + " : Start bundleCreate2")
|
|
||||||
return new Promise(async (resolve, reject) => {
|
return new Promise(async (resolve, reject) => {
|
||||||
|
let lock = null
|
||||||
try {
|
try {
|
||||||
const conn = await MysqlHelpers.createConnection()
|
const conn = await MysqlHelpers.createConnection()
|
||||||
await MysqlHelpers.createTrx(conn)
|
await MysqlHelpers.createTrx(conn)
|
||||||
// console.log("createTrx : " + rltm.device_id)
|
|
||||||
|
|
||||||
let rltmLength = Object.keys(rltm).length
|
let rltmLength = Object.keys(rltm).length
|
||||||
let resLogs = undefined
|
let resLogs = undefined
|
||||||
@ -151,7 +272,6 @@ class GpsTracksModels {
|
|||||||
`INSERT INTO t_gps_tracks SET ?;`,
|
`INSERT INTO t_gps_tracks SET ?;`,
|
||||||
[logs],
|
[logs],
|
||||||
)
|
)
|
||||||
// console.log("insert t_gps_tracks : " + rltm.device_id)
|
|
||||||
|
|
||||||
if (logs.action == "location") {
|
if (logs.action == "location") {
|
||||||
const date = logs.crt_d
|
const date = logs.crt_d
|
||||||
@ -163,9 +283,9 @@ class GpsTracksModels {
|
|||||||
await MysqlHelpers.queryTrx(
|
await MysqlHelpers.queryTrx(
|
||||||
conn,
|
conn,
|
||||||
`
|
`
|
||||||
INSERT INTO tracks_${yy}${mm}
|
INSERT INTO tracks_${yy}${mm}
|
||||||
SET ?
|
SET ?
|
||||||
`,
|
`,
|
||||||
[logs],
|
[logs],
|
||||||
)
|
)
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
@ -175,17 +295,28 @@ class GpsTracksModels {
|
|||||||
JSON.stringify(logs) + `\n`,
|
JSON.stringify(logs) + `\n`,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
// console.log("insert tracks_${yy}${mm} : " + rltm.device_id)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if (rltmLength > 0 && typeof resLogs !== "undefined")
|
if (rltmLength > 0 && typeof resLogs !== "undefined")
|
||||||
rltm.master_id = resLogs.insertId
|
rltm.master_id = resLogs.insertId
|
||||||
|
|
||||||
if (
|
if (
|
||||||
rltmLength > 0 &&
|
rltmLength > 0 &&
|
||||||
rltm.latitude !== null &&
|
rltm.latitude !== null &&
|
||||||
rltm.longitude !== 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(
|
let devices = await MysqlHelpers.queryTrx(
|
||||||
conn,
|
conn,
|
||||||
`SELECT id FROM t_gps_tracks_rltm as rltm WHERE rltm.device_id = ?`,
|
`SELECT id FROM t_gps_tracks_rltm as rltm WHERE rltm.device_id = ?`,
|
||||||
@ -234,7 +365,7 @@ class GpsTracksModels {
|
|||||||
}
|
}
|
||||||
|
|
||||||
await MysqlHelpers.commit(conn)
|
await MysqlHelpers.commit(conn)
|
||||||
// console.log("Commit bundleCreate2 : " + rltm.device_id)
|
|
||||||
resolve({
|
resolve({
|
||||||
type: "success",
|
type: "success",
|
||||||
result: resLogs,
|
result: resLogs,
|
||||||
@ -242,6 +373,16 @@ class GpsTracksModels {
|
|||||||
} catch (err) {
|
} catch (err) {
|
||||||
console.log("ERROR bundleCreate2 : " + rltm.device_id)
|
console.log("ERROR bundleCreate2 : " + rltm.device_id)
|
||||||
reject(err)
|
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)
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@ -34,6 +34,7 @@
|
|||||||
"jsonwebtoken": "^8.5.1",
|
"jsonwebtoken": "^8.5.1",
|
||||||
"moment": "^2.29.1",
|
"moment": "^2.29.1",
|
||||||
"morgan": "^1.10.0",
|
"morgan": "^1.10.0",
|
||||||
|
"mutex": "^1.0.4",
|
||||||
"mysql": "^2.18.1",
|
"mysql": "^2.18.1",
|
||||||
"mysql2": "^3.15.3",
|
"mysql2": "^3.15.3",
|
||||||
"nanoid": "^3.3.1",
|
"nanoid": "^3.3.1",
|
||||||
|
|||||||
Reference in New Issue
Block a user