More work on Postgres support.
Ylian Saint-Hilaire committed
Nov 3, 2021 at 23:42 UTC
11ff538ec5aa1a2fb38d2d8abd6a028722ef3914
1 file changed
+186
-4
db.js
+186
-4
@@ -341,7 +341,10 @@ module.exports.CreateDB = function (parent, func) {
341
if (meshChange) { obj.Set(docs[i]); }
342
}
343
}
344
- if ((obj.databaseType == 4) || (obj.databaseType == 5)) {
344
+ if (obj.databaseType == 6) {
345
+ // Postgres
346
+ sqlDbQuery('DELETE FROM Main WHERE ((extra LIKE (\'mesh/%\')) AND (extra NOT IN ($1)))', [meshlist], func);
347
+ } else if ((obj.databaseType == 4) || (obj.databaseType == 5)) {
348
// MariaDB
349
sqlDbQuery('DELETE FROM Main WHERE (extra LIKE ("mesh/%") AND (extra NOT IN ?)', [meshlist], func);
350
} else if (obj.databaseType == 3) {
@@ -1053,14 +1056,24 @@ module.exports.CreateDB = function (parent, func) {
1056
})
1057
.catch(function (err) { conn.release(); if (func) try { func(err); } catch (ex) { console.log('SQLERR2', ex); } });
1058
}).catch(function (err) { if (func) { try { func(err); } catch (ex) { console.log('SQLERR3', ex); } } });
1056
- } else if ((obj.databaseType == 5) || (obj.databaseType == 6)) { // MySQL or Postgres SQL
1059
+ } else if (obj.databaseType == 5) { // MySQL
1060
Datastore.query(query, args, function (error, results, fields) {
1061
if (error != null) {
1062
if (func) try { func(error); } catch (ex) { console.log('SQLERR4', ex); }
1063
} else {
1064
var docs = [];
1065
for (var i in results) { if (results[i].doc) { docs.push(JSON.parse(results[i].doc)); } }
1063
- //console.log(docs);
1066
+ if (func) { try { func(null, docs); } catch (ex) { console.log('SQLERR5', ex); } }
1067
+ }
1068
+ });
1069
+ } else if (obj.databaseType == 6) { // Postgres SQL
1070
+ Datastore.query(query, args, function (error, results) {
1071
+ if (error != null) {
1072
+ console.log(query, args, error);
1073
+ if (func) try { func(error); } catch (ex) { console.log('SQLERR4', ex); }
1074
+ } else {
1075
+ var docs = [];
1076
+ if (results.command == 'SELECT') { for (var i in results.rows) {if (results.rows[i].doc) { if (typeof results.rows[i].doc == 'string') { docs.push(JSON.parse(results.rows[i].doc)); } else { docs.push(results.rows[i].doc); } } } }
1077
if (func) { try { func(null, docs); } catch (ex) { console.log('SQLERR5', ex); } }
1078
}
1079
});
@@ -1108,7 +1121,176 @@ module.exports.CreateDB = function (parent, func) {
1121
}
1122
1123
function setupFunctions(func) {
1111
- if ((obj.databaseType == 4) || (obj.databaseType == 5) || (obj.databaseType == 6)) {
1124
+ if (obj.databaseType == 6) {
1125
+ // Database actions on the main collection (Postgres)
1126
+ obj.Set = function (value, func) {
1127
+ obj.dbCounters.fileSet++;
1128
+ var extra = null, extraex = null;
1129
+ value = common.escapeLinksFieldNameEx(value);
1130
+ if (value.meshid) { extra = value.meshid; } else if (value.email) { extra = 'email/' + value.email; } else if (value.nodeid) { extra = value.nodeid; }
1131
+ if ((value.type == 'node') && (value.intelamt != null) && (value.intelamt.uuid != null)) { extraex = 'uuid/' + value.intelamt.uuid; }
1132
+ if (value._id == null) { value._id = require('crypto').randomBytes(16).toString('hex'); }
1133
+ sqlDbQuery('INSERT INTO main VALUES ($1, $2, $3, $4, $5, $6) ON CONFLICT (id) DO UPDATE SET type = $2, domain = $3, extra = $4, extraex = $5, doc = $6;', [value._id, (value.type ? value.type : null), ((value.domain != null) ? value.domain : null), extra, extraex, performTypedRecordEncrypt(value)], func);
1134
+ }
1135
+ obj.SetRaw = function (value, func) {
1136
+ obj.dbCounters.fileSet++;
1137
+ var extra = null, extraex = null;
1138
+ if (value.meshid) { extra = value.meshid; } else if (value.email) { extra = 'email/' + value.email; } else if (value.nodeid) { extra = value.nodeid; }
1139
+ if ((value.type == 'node') && (value.intelamt != null) && (value.intelamt.uuid != null)) { extraex = 'uuid/' + value.intelamt.uuid; }
1140
+ if (value._id == null) { value._id = require('crypto').randomBytes(16).toString('hex'); }
1141
+ sqlDbQuery('INSERT INTO main VALUES ($1, $2, $3, $4, $5, $6) ON CONFLICT (id) DO UPDATE SET type = $2, domain = $3, extra = $4, extraex = $5, doc = $6;', [value._id, (value.type ? value.type : null), ((value.domain != null) ? value.domain : null), extra, extraex, performTypedRecordEncrypt(value)], func);
1142
+ }
1143
+ obj.Get = function (_id, func) { sqlDbQuery('SELECT doc FROM main WHERE id = $1', [_id], function (err, docs) { if ((docs != null) && (docs.length > 0) && (docs[0].links != null)) { docs[0] = common.unEscapeLinksFieldName(docs[0]); } func(err, performTypedRecordDecrypt(docs)); }); }
1144
+ obj.GetAll = function (func) { sqlDbQuery('SELECT domain, doc FROM main', null, function (err, docs) { func(err, performTypedRecordDecrypt(docs)); }); }
1145
+ obj.GetHash = function (id, func) { sqlDbQuery('SELECT doc FROM main WHERE id = $1', [id], function (err, docs) { func(err, performTypedRecordDecrypt(docs)); }); }
1146
+ obj.GetAllTypeNoTypeField = function (type, domain, func) { sqlDbQuery('SELECT doc FROM main WHERE type = $1 AND domain = $2', [type, domain], function (err, docs) { if (err == null) { for (var i in docs) { delete docs[i].type } } func(err, performTypedRecordDecrypt(docs)); }); };
1147
+ obj.GetAllTypeNoTypeFieldMeshFiltered = function (meshes, extrasids, domain, type, id, func) {
1148
+ if (id && (id != '')) {
1149
+ sqlDbQuery('SELECT doc FROM main WHERE id = $1 AND type = $2 AND domain = $3 AND extra IN ($4)', [id, type, domain, meshes], function (err, docs) { if (err == null) { for (var i in docs) { delete docs[i].type } } func(err, performTypedRecordDecrypt(docs)); });
1150
+ } else {
1151
+ if (extrasids == null) {
1152
+ sqlDbQuery('SELECT doc FROM main WHERE type = $1 AND domain = $2 AND extra IN ($3)', [type, domain, meshes], function (err, docs) { if (err == null) { for (var i in docs) { delete docs[i].type } } func(err, performTypedRecordDecrypt(docs)); });
1153
+ } else {
1154
+ sqlDbQuery('SELECT doc FROM main WHERE type = $1 AND domain = $2 AND (extra IN ($3) OR id IN ($4))', [type, domain, meshes, extrasids], function (err, docs) { if (err == null) { for (var i in docs) { delete docs[i].type } } func(err, performTypedRecordDecrypt(docs)); });
1155
+ }
1156
+ }
1157
+ };
1158
+ obj.GetAllTypeNodeFiltered = function (nodes, domain, type, id, func) {
1159
+ if (id && (id != '')) {
1160
+ sqlDbQuery('SELECT doc FROM main WHERE id = $1 AND type = $2 AND domain = $3 AND extra IN ($4)', [id, type, domain, nodes], function (err, docs) { if (err == null) { for (var i in docs) { delete docs[i].type } } func(err, performTypedRecordDecrypt(docs)); });
1161
+ } else {
1162
+ sqlDbQuery('SELECT doc FROM main WHERE type = $1 AND domain = $2 AND extra IN ($3)', [type, domain, nodes], function (err, docs) { if (err == null) { for (var i in docs) { delete docs[i].type } } func(err, performTypedRecordDecrypt(docs)); });
1163
+ }
1164
+ };
1165
+ obj.GetAllType = function (type, func) { sqlDbQuery('SELECT doc FROM main WHERE type = $1', [type], function (err, docs) { func(err, performTypedRecordDecrypt(docs)); }); }
1166
+ obj.GetAllIdsOfType = function (ids, domain, type, func) { sqlDbQuery('SELECT doc FROM main WHERE id IN ($1) AND domain = $2 AND type = $3', [ids, domain, type], function (err, docs) { func(err, performTypedRecordDecrypt(docs)); }); }
1167
+ obj.GetUserWithEmail = function (domain, email, func) { sqlDbQuery('SELECT doc FROM main WHERE domain = $1 AND extra = $2', [domain, 'email/' + email], function (err, docs) { func(err, performTypedRecordDecrypt(docs)); }); }
1168
+ obj.GetUserWithVerifiedEmail = function (domain, email, func) { sqlDbQuery('SELECT doc FROM main WHERE domain = $1 AND extra = $2', [domain, 'email/' + email], function (err, docs) { func(err, performTypedRecordDecrypt(docs)); }); }
1169
+ obj.Remove = function (id, func) { sqlDbQuery('DELETE FROM main WHERE id = $1', [id], func); };
1170
+ obj.RemoveAll = function (func) { sqlDbQuery('DELETE FROM main', null, func); };
1171
+ obj.RemoveAllOfType = function (type, func) { sqlDbQuery('DELETE FROM main WHERE type = $1', [type], func); };
1172
+ obj.InsertMany = function (data, func) { var pendingOps = 0; for (var i in data) { pendingOps++; obj.SetRaw(data[i], function () { if (--pendingOps == 0) { func(); } }); } }; // Insert records directly, no link escaping
1173
+ obj.RemoveMeshDocuments = function (id, func) { sqlDbQuery('DELETE FROM main WHERE extra = $1', [id], function () { sqlDbQuery('DELETE FROM main WHERE id = $1', ['nt' + id], func); }); };
1174
+ obj.MakeSiteAdmin = function (username, domain) { obj.Get('user/' + domain + '/' + username, function (err, docs) { if ((err == null) && (docs.length == 1)) { docs[0].siteadmin = 0xFFFFFFFF; obj.Set(docs[0]); } }); };
1175
+ obj.DeleteDomain = function (domain, func) { sqlDbQuery('DELETE FROM main WHERE domain = $1', [domain], func); };
1176
+ obj.SetUser = function (user) { if (user.subscriptions != null) { var u = Clone(user); if (u.subscriptions) { delete u.subscriptions; } obj.Set(u); } else { obj.Set(user); } };
1177
+ obj.dispose = function () { for (var x in obj) { if (obj[x].close) { obj[x].close(); } delete obj[x]; } };
1178
+ obj.getLocalAmtNodes = function (func) { sqlDbQuery('SELECT doc FROM main WHERE (type = \'node\') AND (extraex IS NOT NULL)', null, function (err, docs) { var r = []; if (err == null) { for (var i in docs) { if (docs[i].host != null) { r.push(docs[i]); } } } func(err, r); }); };
1179
+ obj.getAmtUuidMeshNode = function (domainid, mtype, uuid, func) { sqlDbQuery('SELECT doc FROM main WHERE domain = $1 AND extraex = $2', [domainid, 'uuid/' + uuid], func); };
1180
+ obj.isMaxType = function (max, type, domainid, func) { if (max == null) { func(false); } else { sqlDbExec('SELECT COUNT(id) FROM main WHERE domain = $1 AND type = $2', [domainid, type], function (err, response) { func((response['COUNT(id)'] == null) || (response['COUNT(id)'] > max), response['COUNT(id)']) }); } }
1181
+
1182
+ // Database actions on the events collection
1183
+ obj.GetAllEvents = function (func) { sqlDbQuery('SELECT doc FROM events', null, func); };
1184
+ obj.StoreEvent = function (event, func) {
1185
+ obj.dbCounters.eventsSet++;
1186
+ var batchQuery = [['INSERT INTO events VALUES ($1, $2, $3, $4, $5, $6, $7)', [null, event.time, ((typeof event.domain == 'string') ? event.domain : null), event.action, event.nodeid ? event.nodeid : null, event.userid ? event.userid : null, event]]];
1187
+ for (var i in event.ids) { if (event.ids[i] != '*') { batchQuery.push(['INSERT INTO eventids VALUES (LAST_INSERT_ID(), $1)', [event.ids[i]]]); } }
1188
+ sqlDbBatchExec(batchQuery, function (err, docs) { if (func != null) { func(err, docs); } });
1189
+ };
1190
+ obj.GetEvents = function (ids, domain, func) {
1191
+ if (ids.indexOf('*') >= 0) {
1192
+ sqlDbQuery('SELECT doc FROM events WHERE (domain = $1) ORDER BY time DESC', [domain], func);
1193
+ } else {
1194
+ sqlDbQuery('SELECT doc FROM events JOIN eventids ON id = fkid WHERE (domain = $1 AND target IN ($2)) GROUP BY id ORDER BY time DESC', [domain, ids], func);
1195
+ }
1196
+ };
1197
+ obj.GetEventsWithLimit = function (ids, domain, limit, func) {
1198
+ if (ids.indexOf('*') >= 0) {
1199
+ sqlDbQuery('SELECT doc FROM events WHERE (domain = $1) ORDER BY time DESC LIMIT $2', [domain, limit], func);
1200
+ } else {
1201
+ sqlDbQuery('SELECT doc FROM events JOIN eventids ON id = fkid WHERE (domain = $1 AND target IN ($2)) GROUP BY id ORDER BY time DESC LIMIT $3', [domain, ids, limit], func);
1202
+ }
1203
+ };
1204
+ obj.GetUserEvents = function (ids, domain, username, func) {
1205
+ const userid = 'user/' + domain + '/' + username.toLowerCase();
1206
+ if (ids.indexOf('*') >= 0) {
1207
+ sqlDbQuery('SELECT doc FROM events WHERE (domain = $1 AND userid = $2) ORDER BY time DESC', [domain, userid], func);
1208
+ } else {
1209
+ sqlDbQuery('SELECT doc FROM events JOIN eventids ON id = fkid WHERE (domain = $1 AND userid = $2 AND target IN ($3)) GROUP BY id ORDER BY time DESC', [domain, userid, ids], func);
1210
+ }
1211
+ };
1212
+ obj.GetUserEventsWithLimit = function (ids, domain, username, limit, func) {
1213
+ const userid = 'user/' + domain + '/' + username.toLowerCase();
1214
+ if (ids.indexOf('*') >= 0) {
1215
+ sqlDbQuery('SELECT doc FROM events WHERE (domain = $1 AND userid = $2) ORDER BY time DESC LIMIT ?', [domain, userid, limit], func);
1216
+ } else {
1217
+ sqlDbQuery('SELECT doc FROM events JOIN eventids ON id = fkid WHERE (domain = $1 AND userid = $2 AND target IN ($3)) GROUP BY id ORDER BY time DESC LIMIT ?', [domain, userid, ids, limit], func);
1218
+ }
1219
+ };
1220
+ obj.GetEventsTimeRange = function (ids, domain, msgids, start, end, func) {
1221
+ if (ids.indexOf('*') >= 0) {
1222
+ sqlDbQuery('SELECT doc FROM events WHERE ((domain = $1) AND (time BETWEEN $2 AND ?)) ORDER BY time', [domain, start, end], func);
1223
+ } else {
1224
+ sqlDbQuery('SELECT doc FROM events JOIN eventids ON id = fkid WHERE ((domain = $1) AND (target IN ($2)) AND (time BETWEEN $3 AND $4)) GROUP BY id ORDER BY time', [domain, ids, start, end], func);
1225
+ }
1226
+ };
1227
+ //obj.GetUserLoginEvents = function (domain, username, func) { } // TODO
1228
+ obj.GetNodeEventsWithLimit = function (nodeid, domain, limit, func) { sqlDbQuery('SELECT doc FROM events WHERE (nodeid = $1) AND (domain = $2) ORDER BY time DESC LIMIT $3', [nodeid, domain, limit], func); };
1229
+ obj.GetNodeEventsSelfWithLimit = function (nodeid, domain, userid, limit, func) { sqlDbQuery('SELECT doc FROM events WHERE (nodeid = $1) AND (domain = $2) AND ((userid = $3) OR (userid IS NULL)) ORDER BY time DESC LIMIT $4', [nodeid, domain, userid, limit], func); };
1230
+ obj.RemoveAllEvents = function (domain) { sqlDbQuery('DELETE FROM events', null, function (err, docs) { }); };
1231
+ obj.RemoveAllNodeEvents = function (domain, nodeid) { sqlDbQuery('DELETE FROM events WHERE domain = $1 AND nodeid = $2', [domain, nodeid], function (err, docs) { }); };
1232
+ obj.RemoveAllUserEvents = function (domain, userid) { sqlDbQuery('DELETE FROM events WHERE domain = $1 AND userid = $2', [domain, userid], function (err, docs) { }); };
1233
+ obj.GetFailedLoginCount = function (username, domainid, lastlogin, func) { sqlDbExec('SELECT COUNT(id) FROM events WHERE action = "authfail" AND domain = $1 AND userid = $2 AND time > $3', [domainid, 'user/' + domainid + '/' + username.toLowerCase(), lastlogin], function (err, response) { func(err == null ? response['COUNT(id)'] : 0); }); }
1234
+
1235
+ // Database actions on the power collection
1236
+ obj.getAllPower = function (func) { sqlDbQuery('SELECT doc FROM power', null, func); };
1237
+ obj.storePowerEvent = function (event, multiServer, func) { obj.dbCounters.powerSet++; if (multiServer != null) { event.server = multiServer.serverid; } sqlDbQuery('INSERT INTO power VALUES (DEFAULT, $1, $2, $3)', [event.time, event.nodeid ? event.nodeid : null, event], func); };
1238
+ obj.getPowerTimeline = function (nodeid, func) { sqlDbQuery('SELECT doc FROM power WHERE ((nodeid = $1) OR (nodeid = "*")) ORDER BY time ASC', [nodeid], func); };
1239
+ obj.removeAllPowerEvents = function () { sqlDbQuery('DELETE FROM power', null, function (err, docs) { }); };
1240
+ obj.removeAllPowerEventsForNode = function (nodeid) { sqlDbQuery('DELETE FROM power WHERE nodeid = $1', [nodeid], function (err, docs) { }); };
1241
+
1242
+ // Database actions on the SMBIOS collection
1243
+ obj.GetAllSMBIOS = function (func) { sqlDbQuery('SELECT doc FROM smbios', null, func); };
1244
+ obj.SetSMBIOS = function (smbios, func) { var expire = new Date(smbios.time); expire.setMonth(expire.getMonth() + 6); sqlDbQuery('INSERT INTO smbios VALUES ($1, $2, $3, $4) ON CONFLICT (id) DO UPDATE SET time = $2, expire = $3, doc = $4', [smbios._id, smbios.time, expire, smbios], func); };
1245
+ obj.RemoveSMBIOS = function (id) { sqlDbQuery('DELETE FROM smbios WHERE id = $1', [id], function (err, docs) { }); };
1246
+ obj.GetSMBIOS = function (id, func) { sqlDbQuery('SELECT doc FROM smbios WHERE id = $1', [id], func); };
1247
+
1248
+ // Database actions on the Server Stats collection
1249
+ obj.SetServerStats = function (data, func) { sqlDbQuery('INSERT INTO serverstats VALUES ($1, $2, $3) ON CONFLICT (time) DO UPDATE SET expire = $2, doc = $3', [data.time, data.expire, data], func); };
1250
+ obj.GetServerStats = function (hours, func) { var t = new Date(); t.setTime(t.getTime() - (60 * 60 * 1000 * hours)); sqlDbQuery('SELECT doc FROM main WHERE time < $1', [t], func); }; // TODO: Expire old entries
1251
+
1252
+ // Read a configuration file from the database
1253
+ obj.getConfigFile = function (path, func) { obj.Get('cfile/' + path, func); }
1254
+
1255
+ // Write a configuration file to the database
1256
+ obj.setConfigFile = function (path, data, func) { obj.Set({ _id: 'cfile/' + path, type: 'cfile', data: data.toString('base64') }, func); }
1257
+
1258
+ // List all configuration files
1259
+ obj.listConfigFiles = function (func) { sqlDbQuery('SELECT doc FROM main WHERE type = "cfile" ORDER BY id', func); }
1260
+
1261
+ // Get all configuration files
1262
+ obj.getAllConfigFiles = function (password, func) {
1263
+ obj.file.find({ type: 'cfile' }).toArray(function (err, docs) {
1264
+ if (err != null) { func(null); return; }
1265
+ var r = null;
1266
+ for (var i = 0; i < docs.length; i++) {
1267
+ var name = docs[i]._id.split('/')[1];
1268
+ var data = obj.decryptData(password, docs[i].data);
1269
+ if (data != null) { if (r == null) { r = {}; } r[name] = data; }
1270
+ }
1271
+ func(r);
1272
+ });
1273
+ }
1274
+
1275
+ // Get database information (TODO: Complete this)
1276
+ obj.getDbStats = function (func) {
1277
+ obj.stats = { c: 4 };
1278
+ sqlDbExec('SELECT COUNT(id) FROM main', null, function (err, response) { obj.stats.meshcentral = response['COUNT(id)']; if (--obj.stats.c == 0) { delete obj.stats.c; func(obj.stats); } });
1279
+ sqlDbExec('SELECT COUNT(time) FROM serverstats', null, function (err, response) { obj.stats.serverstats = response['COUNT(time)']; if (--obj.stats.c == 0) { delete obj.stats.c; func(obj.stats); } });
1280
+ sqlDbExec('SELECT COUNT(id) FROM power', null, function (err, response) { obj.stats.power = response['COUNT(id)']; if (--obj.stats.c == 0) { delete obj.stats.c; func(obj.stats); } });
1281
+ sqlDbExec('SELECT COUNT(id) FROM smbios', null, function (err, response) { obj.stats.smbios = response['COUNT(id)']; if (--obj.stats.c == 0) { delete obj.stats.c; func(obj.stats); } });
1282
+ }
1283
+
1284
+ // Plugin operations
1285
+ if (obj.pluginsActive) {
1286
+ obj.addPlugin = function (plugin, func) { sqlDbQuery('INSERT INTO plugin VALUES ($1, $2)', [null, value], func); }; // Add a plugin
1287
+ obj.getPlugins = function (func) { sqlDbQuery('SELECT doc FROM plugin', null, func); }; // Get all plugins
1288
+ obj.getPlugin = function (id, func) { sqlDbQuery('SELECT doc FROM plugin WHERE id = $1', [id], func); }; // Get plugin
1289
+ obj.deletePlugin = function (id, func) { sqlDbQuery('DELETE FROM plugin WHERE id = $1', [id], func); }; // Delete plugin
1290
+ obj.setPluginStatus = function (id, status, func) { obj.getPlugin(id, function (err, docs) { if ((err == null) && (docs.length == 1)) { docs[0].status = status; obj.updatePlugin(id, docs[0], func); } }); };
1291
+ obj.updatePlugin = function (id, args, func) { delete args._id; sqlDbQuery('INSERT INTO plugin VALUES ($1, $2) ON CONFLICT (id) DO UPDATE SET doc = $2', [id, args], func); };
1292
+ }
1293
+ } else if ((obj.databaseType == 4) || (obj.databaseType == 5)) {
1294
// Database actions on the main collection (MariaDB or MySQL)
1295
obj.Set = function (value, func) {
1296
obj.dbCounters.fileSet++;