Compare commits

..

26 Commits

Author SHA1 Message Date
b9eede0857 update 2026-07-29 09:27:34 +07:00
319dbbe7bb update 2026-07-29 08:59:49 +07:00
504fb02c3d update 2026-07-29 08:54:53 +07:00
c7215b3caf update 2026-07-28 15:22:36 +07:00
e86e29dab6 update 2026-07-28 15:07:37 +07:00
dd0b95f194 update release connection 2026-07-28 15:06:19 +07:00
856bce46ea update 2026-07-28 14:32:40 +07:00
bb293c9851 update 2026-07-28 14:31:19 +07:00
3522d9db0b update 2026-07-28 13:27:33 +07:00
a6bbf5b6ae UPDATE 2026-07-28 13:24:07 +07:00
f1ea52f9e1 update 2026-07-28 12:54:49 +07:00
8e0ac97945 update 2026-07-28 12:31:10 +07:00
16c3353b21 update: add mutex in create bundle realtime data 2026-07-28 11:57:25 +07:00
548f89059b disable proxy 2026-07-17 18:53:34 +07:00
40b1aa2bc5 update 2026-07-16 13:38:29 +07:00
8322ac1ae1 update 2026-07-16 13:34:55 +07:00
7aa6c72ee4 update 2026-07-16 11:05:30 +07:00
3188098716 update 2026-07-16 10:42:20 +07:00
92b46faaea update 2026-07-16 08:54:25 +07:00
1d525074a1 update 2026-07-16 08:52:20 +07:00
332692e01c update 2026-07-14 16:51:51 +07:00
6e5eaf8268 update 2026-07-14 16:29:09 +07:00
b95d0b2bc5 update 2026-07-14 10:51:34 +07:00
9ab77e6db4 update 2026-07-13 16:48:30 +07:00
dd8dbfa3ff update 2026-07-13 13:22:33 +07:00
a0cd3a7d8b update 2026-07-13 10:54:43 +07:00
9 changed files with 532 additions and 217 deletions

106
app.js
View File

@ -1,3 +1,5 @@
//app.js
require("dotenv").config({ path: ".env" })
require("events").EventEmitter.prototype._maxListeners = 30
const morgan = require("morgan")
@ -14,6 +16,7 @@ const LibMail = require("./library/LibMail")
const LibHelper = require("./library/LibHelper")
const nanoid = require("nanoid").nanoid
const LibWinston = require("./library/LibWinston")
const cron = require("node-cron")
const express = require("express")
// const routes = require("./routes/routes")
@ -23,6 +26,8 @@ const app = express()
const logName = "libUdp"
const Logger = LibWinston.initialize(logName)
process.env["NODE_TLS_REJECT_UNAUTHORIZED"] = 0;
async function commitMessage(now, logDevice) {
try {
if (!logDevice.original_hex) {
@ -288,6 +293,102 @@ async function commitMessage(now, logDevice) {
}
}
async function checkMaintenanceMileage(vhc) {
try {
if (!vhc || !vhc.id) return false
const lastHistory = await db.query(
`SELECT odometer FROM t_vehicles_maintenance_history
WHERE vhc_id = ? AND isdeleted = 0
ORDER BY id DESC LIMIT 1`,
[vhc.id],
)
if (!lastHistory?.length) return false
const lastOdometer = Number(lastHistory[0].odometer || 0)
const currentMileage = Number(vhc.sum_milleage || 0)
const diff = currentMileage - lastOdometer
const thresholds = [5000, 10000, 150000]
if (!thresholds.includes(diff)) return false
await db.query(
`INSERT INTO t_vehicles_maintenance_history SET vhc_id = ?, odometer = ?`,
[vhc.id, currentMileage],
)
const users = await db.query(
`SELECT phone FROM t_users WHERE status_sms = 1 AND FIND_IN_SET(?, vhc) > 0`,
[vhc.id],
)
if (!users?.length) return false
const nopol = (vhc.nopol1 ?? "") + (vhc.nopol2 ?? "") + (vhc.nopol3 ?? "")
const rawMsg = `${nopol} | MAINTENANCE | ${diff} km sejak maintenance terakhir`.substring(0, 150).trim()
const countryCode = process.env.SMS_COUNTRY_CODE ?? "670"
const smsHost = process.env.SMS_HOST ?? "http://192.168.40.2:8181/submitsm/ka"
await Promise.allSettled(
users
.filter((u) => u.phone)
.map(async (user) => {
try {
await axios.get(`${smsHost}`, {
params: { to: `${countryCode}${user.phone}`, msg: rawMsg, sender: "Maintenance FROTA" },
})
} catch (smsErr) {
console.error("Error sending SMS:", smsErr.response ? JSON.stringify(smsErr.response.data) : smsErr.message)
}
try {
await db.query(`INSERT INTO sms_notif SET to_no = ?, message = ?`, [user.phone, rawMsg])
} catch (dbErr) {
console.error("Error logging SMS to DB:", dbErr.message)
}
}),
)
} catch (e) {
console.error("Error checkMaintenanceMileage:", e)
}
}
async function checkMaintenanceMileageScheduler() {
console.log("Maintenance mileage check job executed:", moment().format("YYYY-MM-DD HH:mm:ss"))
try {
console.time("Maintenance mileage check")
// ambil semua vehicle aktif beserta mileage & nopol
// vhc.id di sini didapat dari kolom `id` tabel t_vehicles, BUKAN dari device
const vehicles = await db.query(
`SELECT id, nopol1, nopol2, nopol3, sum_milleage
FROM t_vehicles
WHERE dlt IS NULL`,
[],
)
if (!vehicles?.length) {
console.log("No active vehicles found.")
console.timeEnd("Maintenance mileage check")
return
}
for (const vhc of vehicles) {
try {
await checkMaintenanceMileage(vhc)
} catch (err) {
console.error(`Error checking maintenance for vhc_id ${vhc.id}:`, err.message)
}
}
console.timeEnd("Maintenance mileage check")
} catch (error) {
console.error("Error in checkMaintenanceMileageScheduler:", error.message)
}
console.log("Maintenance mileage check job completed.")
}
const devices = []
const netConn = require("./config/netConn")
/**
@ -519,8 +620,8 @@ net.createServer(
commitMessage(now, logDevice)
})
c.on("end", () => {})
c.on("close", () => {})
c.on("end", () => { })
c.on("close", () => { })
c.on("drain", (a) => {
console.log("client drain", a)
@ -643,3 +744,4 @@ udp.on("listening", () => {
})
udp.bind(process.env.PORT_UDP)
cron.schedule("0 * * * *", () => { checkMaintenanceMileageScheduler() })

View File

@ -3,7 +3,7 @@ const request = {
urlBase: 'https://nominatim.openstreetmap.org',
urlPath: 'reverse',
// urlFull: 'https://nominatim.openstreetmap.org/reverse',
urlFull: 'https://brilianapps.britimorleste.tl/nominatim/reverse',
urlFull: 'http://api-gateway.corp.timortelecom.tl/nominatim/reverse',
method: 'GET',
data: {
format: 'geojson' // xml,json,jsonv2,geojson(prefer),geocodejson

View File

@ -69,7 +69,8 @@ const go = async () => {
fulladdress: encodeURIComponent(decodeURIComponent(sameAddr[0].fulladdress)),
type_reverse_geo: sameAddr[0].type_reverse_geo,
stts_reverse_geo: GpsTracksModels.STTS_REVERSE_GEO_SC,
log_reverse_geo: sameAddr[0].log_reverse_geo,
// log_reverse_geo: sameAddr[0].log_reverse_geo,
log_reverse_geo: "",
crt: now,
crt_format: moment.unix(now).format("YYYY-MM-DD HH:mm:ss"),
}
@ -82,10 +83,10 @@ const go = async () => {
const resp = await axInstance.get(`${urlBase}?${params.toString()}`, {
headers: { "User-Agent": `movana-fleet-management-` + i },
proxy: false,
})
const respData = resp.data || {}
if (resp.status === 200) {
console.log("SUCCESS respReverseGeo:", track.id)
if (respData.features.length < 1) {
GpsTracksModels.create2Address({
device_id: track.device_id,
@ -93,7 +94,8 @@ const go = async () => {
lat: track.latitude,
lng: track.longitude,
stts_reverse_geo: GpsTracksModels.STTS_REVERSE_GEO_LOST,
log_reverse_geo: JSON.stringify(respData),
// log_reverse_geo: JSON.stringify(respData),
log_reverse_geo: "",
crt: now,
crt_format: moment.unix(now).format("YYYY-MM-DD HH:mm:ss"),
})
@ -118,9 +120,10 @@ const go = async () => {
postcode: respAddr.postcode,
fulladdress: encodeURIComponent(respData.features[0].properties.display_name),
stts_reverse_geo: GpsTracksModels.STTS_REVERSE_GEO_SC,
log_reverse_geo: respData.features[0].properties
? JSON.stringify(respData.features[0].properties)
: null,
// log_reverse_geo: respData.features[0].properties
// ? JSON.stringify(respData.features[0].properties)
// : null,
log_reverse_geo: "",
crt: now,
crt_format: moment.unix(now).format("YYYY-MM-DD HH:mm:ss"),
}
@ -254,7 +257,8 @@ const go = async () => {
lat: track.latitude,
lng: track.longitude,
stts_reverse_geo: GpsTracksModels.STTS_REVERSE_GEO_ER,
log_reverse_geo: JSON.stringify(respData),
// log_reverse_geo: JSON.stringify(respData),
log_reverse_geo: "",
crt: now,
crt_format: moment.unix(now).format("YYYY-MM-DD HH:mm:ss"),
})
@ -282,7 +286,7 @@ const go = async () => {
lng: track.longitude,
stts_reverse_geo: GpsTracksModels.STTS_REVERSE_GEO_ER,
// log_reverse_geo: JSON.stringify(respData.data),
log_reverse_geo: stringify(respData.data),
log_reverse_geo: "",
crt: now,
crt_format: moment.unix(now).format("YYYY-MM-DD HH:mm:ss"),
})

View File

@ -181,189 +181,91 @@ async function tripGrouping() {
const endOfMonth = histMonth.clone().endOf("month").unix() - TIMEFIX
const q2 = `
insert into trips
(id,name,nopol1,vhc_id,mileage,start,finish,startMileage,finishMileage,startLoc,finishLoc,pool_code,dc_code,fuel_consume,row_count)
WITH
sites AS (
SELECT 'Bitterlaun' AS name, 'Hera' AS village, 'Cristo Rei' AS subdistrict, 'Dili' AS district, -8.548328 AS lat, 125.623808 AS lng
UNION ALL SELECT 'Bebonuk', 'Comoro', 'Dom Aleixo', 'Dili', -8.547527, 125.548280
UNION ALL SELECT 'Balide', 'Santa Cruz', 'Nain Feto', 'Dili', -8.564396, 125.582151
)
, gaps AS (
SELECT
CASE
WHEN (crt_d - LAG(crt_d, 1, NULL) OVER (PARTITION BY vhc_id ORDER BY crt_d)) > 3600
THEN 1 ELSE 0
END AS isStop,
(
SELECT CONCAT('Site ', s.name, ', ', s.village, ', ', s.subdistrict, ', ', s.district, ', Timor-Leste')
FROM sites s
WHERE ST_Distance_Sphere(POINT(t.longitude, t.latitude), POINT(s.lng, s.lat)) <= 500
ORDER BY ST_Distance_Sphere(POINT(t.longitude, t.latitude), POINT(s.lng, s.lat)) ASC
LIMIT 1
) AS site_address,
t.*
FROM ${histTableName} t
WHERE
t.latitude IS NOT NULL
AND t.longitude IS NOT NULL
AND t.action = 'location'
AND t.crt_d BETWEEN ? AND ?
)
, trips AS (
SELECT
CASE
WHEN ignition = 4
AND LAG(ignition, 1, 0) OVER (PARTITION BY vhc_id ORDER BY crt_d) <> 4
or LAG(isStop, 1, 0) over (PARTITION BY vhc_id ORDER BY crt_d) = 1
THEN 1 ELSE 0
END AS trip_start,
g.*
FROM gaps g
)
, numbered AS (
SELECT
*,
SUM(trip_start) OVER (PARTITION BY vhc_id ORDER BY crt_d) AS trip_id
FROM trips
where
ignition = 4
and isStop = 0
),
agg AS (
SELECT
v.id,
v.name,
v.nopol1,
vhc_id,
ROW_NUMBER() OVER (PARTITION BY v.id ORDER BY MIN(a.crt_d)) AS trip_id,
max(a.vhc_milleage) - min(a.vhc_milleage) AS mileage,
MIN(a.crt_d) AS start,
MAX(a.crt_d) AS finish,
MIN(a.vhc_milleage) AS startMileage,
MAX(a.vhc_milleage) AS finishMileage,
(
SELECT COALESCE(n.site_address, (SELECT fulladdress FROM t_gps_tracks_address WHERE master_id = n.id LIMIT 1))
FROM numbered n
WHERE n.id = MIN(a.id)
LIMIT 1
) AS startLoc,
(
SELECT COALESCE(n.site_address, (SELECT fulladdress FROM t_gps_tracks_address WHERE master_id = n.id LIMIT 1))
FROM numbered n
WHERE n.id = MAX(a.id)
LIMIT 1
) AS finishLoc,
COUNT(*) AS row_count,
max(fuel_count) - min(fuel_count) AS fuel_consume
FROM t_vehicles v
LEFT JOIN numbered a ON a.vhc_id = v.id
WHERE
v.dlt is null and trip_id != 0
GROUP BY v.id, a.trip_id
HAVING COUNT(*) > 1
)
SELECT
agg.id,name,nopol1,vhc_id,mileage,start,finish,startMileage,finishMileage,startLoc,finishLoc,
tvd.pool_code, tvd.dc_code,fuel_consume,
row_count
FROM agg agg
join t_vehicles_detail tvd on tvd.vid = agg.id
ORDER BY agg.id, trip_id
ON DUPLICATE KEY UPDATE
mileage = values(mileage),
start = values(start),
finish = values(finish),
startMileage = values(startMileage),
finishMileage = values(finishMileage),
startLoc = values(startLoc),
finishLoc = values(finishLoc),
row_count = values(row_count),
fuel_consume = values(fuel_consume)
`
// const q2 = `
// insert into trips
// (id,name,nopol1,vhc_id,mileage,start,finish,startMileage,finishMileage,startLoc,finishLoc,pool_code,dc_code,fuel_consume,row_count)
// WITH
// gaps AS (
// SELECT
// -- previous gap since previous row > 1 hour (3600s)
// CASE
// WHEN (crt_d - LAG(crt_d, 1, NULL) OVER (PARTITION BY vhc_id ORDER BY crt_d)) > 3600
// THEN 1 ELSE 0
// END AS isStop,
// t.*
// FROM ${histTableName} t
// WHERE
// t.latitude IS NOT NULL
// AND t.longitude IS NOT NULL
// AND t.action = 'location'
// AND t.crt_d BETWEEN ? AND ?
// )
// , trips AS (
// SELECT
// -- mark the start of a trip when ignition=4 and previous ignition <> 4
// CASE
// WHEN ignition = 4
// AND LAG(ignition, 1, 0) OVER (PARTITION BY vhc_id ORDER BY crt_d) <> 4
// or LAG(isStop, 1, 0) over (PARTITION BY vhc_id ORDER BY crt_d) = 1
// THEN 1 ELSE 0
// END AS trip_start,
// g.*
// FROM gaps g
// )
// , numbered AS (
// SELECT
// *,
// -- assign a trip_id by cumulative sum of trip_start
// SUM(trip_start) OVER (PARTITION BY vhc_id ORDER BY crt_d) AS trip_id
// FROM trips
// where
// ignition = 4
// and isStop = 0
// ),
// agg AS (
// SELECT
// v.id,
// v.name,
// v.nopol1,
// vhc_id,
// ROW_NUMBER() OVER (PARTITION BY v.id ORDER BY MIN(a.crt_d)) AS trip_id,
// max(a.vhc_milleage) - min(a.vhc_milleage) AS mileage,
// MIN(a.crt_d) AS start,
// MAX(a.crt_d) AS finish,
// MIN(a.vhc_milleage) AS startMileage,
// MAX(a.vhc_milleage) AS finishMileage,
// (SELECT fulladdress FROM t_gps_tracks_address WHERE master_id = MIN(a.id) LIMIT 1) AS startLoc,
// (SELECT fulladdress FROM t_gps_tracks_address WHERE master_id = MAX(a.id) LIMIT 1) AS finishLoc,
// COUNT(*) AS row_count,
// max(fuel_count) - min(fuel_count) AS fuel_consume
// FROM t_vehicles v
// LEFT JOIN numbered a ON a.vhc_id = v.id
// WHERE
// v.dlt is null and trip_id != 0
// GROUP BY v.id, a.trip_id
// HAVING COUNT(*) > 1
// )
// SELECT
// agg.id,name,nopol1,vhc_id,mileage,start,finish,startMileage,finishMileage,startLoc,finishLoc,
// tvd.pool_code, tvd.dc_code,fuel_consume,
// row_count
// FROM agg agg
// join t_vehicles_detail tvd on tvd.vid = agg.id
// ORDER BY agg.id, trip_id
// ON DUPLICATE KEY UPDATE
// mileage = values(mileage),
// start = values(start),
// finish = values(finish),
// startMileage = values(startMileage),
// finishMileage = values(finishMileage),
// startLoc = values(startLoc),
// finishLoc = values(finishLoc),
// row_count = values(row_count),
// fuel_consume = values(fuel_consume)
// `
insert into trips
(id,name,nopol1,vhc_id,mileage,start,finish,startMileage,finishMileage,startLoc,finishLoc,pool_code,dc_code,fuel_consume,row_count,start_lat_long,finish_lat_long)
WITH
gaps AS (
SELECT
-- previous gap since previous row > 1 hour (3600s)
CASE
WHEN (crt_d - LAG(crt_d, 1, NULL) OVER (PARTITION BY vhc_id ORDER BY crt_d)) > 3600
THEN 1 ELSE 0
END AS isStop,
t.*
FROM ${histTableName} t
WHERE
t.latitude IS NOT NULL
AND t.longitude IS NOT NULL
AND t.action = 'location'
AND t.crt_d BETWEEN ? AND ?
)
, trips AS (
SELECT
-- mark the start of a trip when ignition=4 and previous ignition <> 4
CASE
WHEN ignition = 4
AND LAG(ignition, 1, 0) OVER (PARTITION BY vhc_id ORDER BY crt_d) <> 4
or LAG(isStop, 1, 0) over (PARTITION BY vhc_id ORDER BY crt_d) = 1
THEN 1 ELSE 0
END AS trip_start,
g.*
FROM gaps g
)
, numbered AS (
SELECT
*,
-- assign a trip_id by cumulative sum of trip_start
SUM(trip_start) OVER (PARTITION BY vhc_id ORDER BY crt_d) AS trip_id
FROM trips
where
ignition = 4
and isStop = 0
),
agg AS (
SELECT
v.id,
v.name,
v.nopol1,
vhc_id,
ROW_NUMBER() OVER (PARTITION BY v.id ORDER BY MIN(a.crt_d)) AS trip_id,
max(a.vhc_milleage) - min(a.vhc_milleage) AS mileage,
MIN(a.crt_d) AS start,
MAX(a.crt_d) AS finish,
MIN(a.vhc_milleage) AS startMileage,
MAX(a.vhc_milleage) AS finishMileage,
(SELECT fulladdress FROM t_gps_tracks_address WHERE master_id = MIN(a.id) LIMIT 1) AS startLoc,
(SELECT fulladdress FROM t_gps_tracks_address WHERE master_id = MAX(a.id) LIMIT 1) AS finishLoc,
(SELECT CONCAT (latitude,', ',longitude) FROM ${histTableName} WHERE id = MIN(a.id) LIMIT 1) AS start_lat_long,
(SELECT CONCAT (latitude,', ',longitude) FROM ${histTableName} WHERE ID = MAX(a.id) LIMIT 1) AS finish_lat_long,
COUNT(*) AS row_count,
max(fuel_count) - min(fuel_count) AS fuel_consume
FROM t_vehicles v
LEFT JOIN numbered a ON a.vhc_id = v.id
WHERE
v.dlt is null and trip_id != 0
GROUP BY v.id, a.trip_id
HAVING COUNT(*) > 1
)
SELECT
agg.id,name,nopol1,vhc_id,mileage,start,finish,startMileage,finishMileage,startLoc,finishLoc,
tvd.pool_code, tvd.dc_code,fuel_consume,
row_count,start_lat_long,finish_lat_long
FROM agg agg
join t_vehicles_detail tvd on tvd.vid = agg.id
ORDER BY agg.id, trip_id
ON DUPLICATE KEY UPDATE
mileage = values(mileage),
start = values(start),
finish = values(finish),
startMileage = values(startMileage),
finishMileage = values(finishMileage),
startLoc = values(startLoc),
finishLoc = values(finishLoc),
row_count = values(row_count),
fuel_consume = values(fuel_consume),
start_lat_long = values(start_lat_long),
finish_lat_long = values(finish_lat_long)
`
const d2 = [startOfMonth, endOfMonth]
const r2 = await db.query(q2, d2)
console.log(`Inserted ${r2.affectedRows} rows into 'trips' table from '${histTableName}'`)

View File

@ -42,7 +42,8 @@ class LibCurl {
country_text: (respAddr.country) ? respAddr.country.toUpperCase() : respAddr.country || null,
postcode: respAddr.postcode,
fulladdress: encodeURIComponent(respData.features[0].properties.display_name),
log_reverse_geo: (respData.features[0].properties) ? JSON.stringify(respData.features[0].properties) : null,
// log_reverse_geo: (respData.features[0].properties) ? JSON.stringify(respData.features[0].properties) : null,
log_reverse_geo:""
};
// if (respAddr.state || respAddr.city) {
addrData.state_id = null;

View File

@ -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)
})
})
}
}

View File

@ -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
@ -246,6 +256,285 @@ 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
// let conn = null
// try {
// 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],
// )
// 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
// await MysqlHelpers.queryTrx(
// conn,
// `
// INSERT INTO tracks_${yy}${mm}
// SET ?
// `,
// [logs]
// )
// }
// } finally {
// if (trackLock) {
// await rltmMutex.unlock(trackLock)
// }
// // conn.release();
// }
// }
// 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: 90000, // maksimal lock dipegang 5 detik
// maxWait: 90000, // 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)
// }
// }
// conn.release(); // release connection
// }
// })
// }
static async get2() {
return new Promise((resolve, reject) => {
let params = []

View File

@ -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",

View File

@ -49,7 +49,8 @@ module.exports = async (job) => {
fulladdress: encodeURIComponent(decodeURIComponent(sameAddr[0].fulladdress)),
type_reverse_geo: sameAddr[0].type_reverse_geo,
stts_reverse_geo: GpsTracksModels.STTS_REVERSE_GEO_SC,
log_reverse_geo: sameAddr[0].log_reverse_geo,
// log_reverse_geo: sameAddr[0].log_reverse_geo,
log_reverse_geo:"",
crt: now,
crt_format: moment.unix(now).format('YYYY-MM-DD HH:mm:ss'),
};
@ -85,7 +86,8 @@ module.exports = async (job) => {
lat: tracks[i].latitude,
lng: tracks[i].longitude,
stts_reverse_geo: GpsTracksModels.STTS_REVERSE_GEO_LOST,
log_reverse_geo: JSON.stringify(respData),
// log_reverse_geo: JSON.stringify(respData),
log_reverse_geo: "",
crt: now,
crt_format: moment.unix(now).format('YYYY-MM-DD HH:mm:ss'),
});
@ -108,7 +110,8 @@ module.exports = async (job) => {
postcode: respAddr.postcode,
fulladdress: encodeURIComponent(respData.features[0].properties.display_name),
stts_reverse_geo: GpsTracksModels.STTS_REVERSE_GEO_SC,
log_reverse_geo: (respData.features[0].properties) ? JSON.stringify(respData.features[0].properties) : null,
// log_reverse_geo: (respData.features[0].properties) ? JSON.stringify(respData.features[0].properties) : null,
log_reverse_geo: "",
crt: now,
crt_format: moment.unix(now).format('YYYY-MM-DD HH:mm:ss'),
};
@ -201,7 +204,8 @@ module.exports = async (job) => {
lat: tracks[i].latitude,
lng: tracks[i].longitude,
stts_reverse_geo: GpsTracksModels.STTS_REVERSE_GEO_ER,
log_reverse_geo: JSON.stringify(respData),
// log_reverse_geo: JSON.stringify(respData),
log_reverse_geo:"",
crt: now,
crt_format: moment.unix(now).format('YYYY-MM-DD HH:mm:ss'),
});
@ -227,7 +231,8 @@ module.exports = async (job) => {
lat: tracks[i].latitude,
lng: tracks[i].longitude,
stts_reverse_geo: GpsTracksModels.STTS_REVERSE_GEO_ER,
log_reverse_geo: JSON.stringify(respData),
// log_reverse_geo: JSON.stringify(respData),
log_reverse_geo:"",
crt: now,
crt_format: moment.unix(now).format('YYYY-MM-DD HH:mm:ss'),
});