Update database files and improvements

- moved base table struct to base_table.json
- moved cloud data struct to data_structure.json
- improvements to database module
- use database.DB from module.export cache instead of global mapping
This commit is contained in:
sairajzero 2022-12-17 01:10:34 +05:30
parent ee8089f858
commit 43a3b16bbe
6 changed files with 535 additions and 512 deletions

24
src/base_tables.json Normal file
View File

@ -0,0 +1,24 @@
{
"LastTxs": {
"ID": "CHAR(34) NOT NULL",
"N": "INT NOT NULL",
"PRIMARY": "KEY (ID)"
},
"Configs": {
"NAME": "VARCHAR(64) NOT NULL",
"VAL": "VARCHAR(512) NOT NULL",
"PRIMARY": "KEY (NAME)"
},
"SuperNodes": {
"FLO_ID": "CHAR(34) NOT NULL",
"PUB_KEY": "CHAR(66) NOT NULL",
"URI": "VARCHAR(256) NOT NULL",
"PRIMARY": "KEY (FLO_ID)"
},
"Applications": {
"APP_NAME": "VARCHAR(64) NOT NULL",
"ADMIN_ID": "CHAR(34) NOT NULL",
"SUB_ADMINS": "VARCHAR(3500)",
"PRIMARY": "KEY (APP_NAME)"
}
}

View File

@ -1,4 +1,5 @@
var DB, _list; //container for database and _list (stored n serving) const { DB } = require("./database");
const { _list } = require("./intra");
global.INVALID = function (message) { global.INVALID = function (message) {
if (!(this instanceof INVALID)) if (!(this instanceof INVALID))
@ -272,11 +273,5 @@ module.exports = {
checkIfRequestSatisfy, checkIfRequestSatisfy,
processRequestFromUser, processRequestFromUser,
processIncomingData, processIncomingData,
processStatusFromUser, processStatusFromUser
set DB(db) {
DB = db;
},
set _list(list) {
_list = list;
}
}; };

33
src/data_structure.json Normal file
View File

@ -0,0 +1,33 @@
{
"H_struct": {
"VECTOR_CLOCK": "vectorClock",
"SENDER_ID": "senderID",
"RECEIVER_ID": "receiverID",
"TYPE": "type",
"APPLICATION": "application",
"TIME": "time",
"PUB_KEY": "pubKey"
},
"B_struct": {
"MESSAGE": "message",
"SIGNATURE": "sign",
"COMMENT": "comment"
},
"L_struct": {
"PROXY_ID": "proxyID",
"STATUS": "status_n",
"LOG_TIME": "log_time"
},
"T_struct": {
"TAG": "tag",
"TAG_TIME": "tag_time",
"TAG_KEY": "tag_key",
"TAG_SIGN": "tag_sign"
},
"F_struct": {
"NOTE": "note",
"NOTE_TIME": "note_time",
"NOTE_KEY": "note_key",
"NOTE_SIGN": "note_sign"
}
}

View File

@ -1,157 +1,173 @@
'use strict'; 'use strict';
var mysql = require('mysql'); var mysql = require('mysql');
const Base_Tables = { const Base_Tables = require("./base_tables.json");
LastTxs: { const { H_struct, B_struct, L_struct, T_struct, F_struct } = require("./data_structure.json");
ID: "CHAR(34) NOT NULL",
N: "INT NOT NULL", var pool; //container for pool
PRIMARY: "KEY (ID)"
}, function initConnection(user, password, dbname, host = 'localhost') {
Configs: { return new Promise((resolve, reject) => {
NAME: "VARCHAR(64) NOT NULL", pool = mysql.createPool({
VAL: "VARCHAR(512) NOT NULL", host: host,
PRIMARY: "KEY (NAME)" user: user,
}, password: password,
SuperNodes: { database: dbname
FLO_ID: "CHAR(34) NOT NULL", });
PUB_KEY: "CHAR(66) NOT NULL", getConnection().then(conn => {
URI: "VARCHAR(256) NOT NULL", conn.release();
PRIMARY: "KEY (FLO_ID)" resolve(DB);
}, }).catch(error => reject(error));
Applications: { });
APP_NAME: "VARCHAR(64) NOT NULL",
ADMIN_ID: "CHAR(34) NOT NULL",
SUB_ADMINS: "VARCHAR(3500)",
PRIMARY: "KEY (APP_NAME)"
}
}; };
const H_struct = { const getConnection = () => new Promise((resolve, reject) => {
VECTOR_CLOCK: "vectorClock", pool.getConnection((error, conn) => {
SENDER_ID: "senderID", if (error)
RECEIVER_ID: "receiverID", reject(error);
TYPE: "type", else
APPLICATION: "application", resolve(conn);
TIME: "time", });
PUB_KEY: "pubKey" });
};
const B_struct = { function queryResolve(sql, values) {
MESSAGE: "message", return new Promise((resolve, reject) => {
SIGNATURE: "sign", getConnection().then(conn => {
COMMENT: "comment" const fn = (err, res) => {
}; conn.release();
(err ? reject(err) : resolve(res));
const L_struct = { };
PROXY_ID: "proxyID", if (values)
STATUS: "status_n", conn.query(sql, values, fn);
LOG_TIME: "log_time" else
}; conn.query(sql, fn);
}).catch(error => reject(error));
const T_struct = { })
TAG: "tag",
TAG_TIME: "tag_time",
TAG_KEY: "tag_key",
TAG_SIGN: "tag_sign"
};
const F_struct = {
NOTE: "note",
NOTE_TIME: "note_time",
NOTE_KEY: "note_key",
NOTE_SIGN: "note_sign",
} }
function Database(user, password, dbname, host = 'localhost') { function queryStream(sql, value, callback) {
const db = {}; return new Promise((resolve, reject) => {
getConnection().then(conn => {
if (!callback && value instanceof Function) {
callback = value;
value = undefined;
}
var err_flag, row_count = 0;
(value ? conn.query(sql, values) : conn.query(sql))
.on('error', err => err_flag = err)
.on('result', row => {
row_count++;
conn.pause();
callback(row);
conn.resume();
}).on('end', _ => {
conn.release();
err_flag ? reject(err_flag) : resolve(row_count);
});
})
})
}
db.createBase = function () { const DB = {};
/* Base Tables */
DB.createBase = function () {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statements = []; let statements = [];
for (let t in Base_Tables) for (let t in Base_Tables)
statements.push("CREATE TABLE IF NOT EXISTS " + t + "( " + statements.push("CREATE TABLE IF NOT EXISTS " + t + "( " +
Object.keys(Base_Tables[t]).map(a => a + " " + Base_Tables[t][a]).join(", ") + " )"); Object.keys(Base_Tables[t]).map(a => a + " " + Base_Tables[t][a]).join(", ") + " )");
Promise.all(statements.map(s => db.query(s))) Promise.all(statements.map(s => queryResolve(s)))
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.setLastTx = function (id, n) { DB.setLastTx = function (id, n) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "INSERT INTO LastTxs (ID, N) VALUES (?, ?)" + let statement = "INSERT INTO LastTxs (ID, N) VALUES (?, ?)" +
" ON DUPLICATE KEY UPDATE N=?"; " ON DUPLICATE KEY UPDATE N=?";
db.query(statement, [id, n, n]) queryResolve(statement, [id, n, n])
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.setConfig = function (name, value) { DB.setConfig = function (name, value) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "INSERT INTO Configs (NAME, VAL) VALUES (?, ?)" + let statement = "INSERT INTO Configs (NAME, VAL) VALUES (?, ?)" +
" ON DUPLICATE KEY UPDATE VAL=?"; " ON DUPLICATE KEY UPDATE VAL=?";
db.query(statement, [name, value, value]) queryResolve(statement, [name, value, value])
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.addSuperNode = function (id, pubKey, uri) { DB.addSuperNode = function (id, pubKey, uri) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "INSERT INTO SuperNodes (FLO_ID, PUB_KEY, URI) VALUES (?, ?, ?)" + let statement = "INSERT INTO SuperNodes (FLO_ID, PUB_KEY, URI) VALUES (?, ?, ?)" +
" ON DUPLICATE KEY UPDATE URI=?"; " ON DUPLICATE KEY UPDATE URI=?";
db.query(statement, [id, pubKey, uri, uri]) queryResolve(statement, [id, pubKey, uri, uri])
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.rmSuperNode = function (id) { DB.rmSuperNode = function (id) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "DELETE FROM SuperNodes" + let statement = "DELETE FROM SuperNodes" +
" WHERE FLO_ID=?"; " WHERE FLO_ID=?";
db.query(statement, id) queryResolve(statement, id)
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.setSubAdmin = function (appName, subAdmins) { DB.updateSuperNode = function (id, uri) {
return new Promise((resolve, reject) => {
let statement = "UPDATE SuperNodes SET URI=? WHERE FLO_ID=?";
queryResolve(statement, [uri, id])
.then(result => resolve(result))
.catch(error => reject(error));
})
}
DB.setSubAdmin = function (appName, subAdmins) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "UPDATE Applications" + let statement = "UPDATE Applications" +
" SET SUB_ADMINS=?" + " SET SUB_ADMINS=?" +
" WHERE APP_NAME=?"; " WHERE APP_NAME=?";
db.query(statement, [subAdmins.join(","), appName]) queryResolve(statement, [subAdmins.join(","), appName])
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.addApp = function (appName, adminID) { DB.addApp = function (appName, adminID) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "INSERT INTO Applications (APP_NAME, ADMIN_ID) VALUES (?, ?)" + let statement = "INSERT INTO Applications (APP_NAME, ADMIN_ID) VALUES (?, ?)" +
" ON DUPLICATE KEY UPDATE ADMIN_ID=?"; " ON DUPLICATE KEY UPDATE ADMIN_ID=?";
db.query(statement, [appName, adminID, adminID]) queryResolve(statement, [appName, adminID, adminID])
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.rmApp = function (appName) { DB.rmApp = function (appName) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "DELETE FROM Applications" + let statement = "DELETE FROM Applications" +
" WHERE APP_NAME=" + appName; " WHERE APP_NAME=" + appName;
db.query(statement) queryResolve(statement)
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.getBase = function () { DB.getBase = function () {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let tables = Object.keys(Base_Tables); let tables = Object.keys(Base_Tables);
Promise.all(tables.map(t => db.query("SELECT * FROM " + t))).then(result => { Promise.all(tables.map(t => queryResolve("SELECT * FROM " + t))).then(result => {
let tmp = Object.fromEntries(tables.map((t, i) => [t, result[i]])); let tmp = Object.fromEntries(tables.map((t, i) => [t, result[i]]));
result = {}; result = {};
result.lastTx = Object.fromEntries(tmp.LastTxs.map(a => [a.ID, a.N])); result.lastTx = Object.fromEntries(tmp.LastTxs.map(a => [a.ID, a.N]));
@ -165,9 +181,12 @@ function Database(user, password, dbname, host = 'localhost') {
resolve(result); resolve(result);
}).catch(error => reject(error)); }).catch(error => reject(error));
}); });
}; };
db.createTable = function (snID) {
/* Supernode Tables */
DB.createTable = function (snID) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "CREATE TABLE IF NOT EXISTS _" + snID + " ( " + let statement = "CREATE TABLE IF NOT EXISTS _" + snID + " ( " +
H_struct.VECTOR_CLOCK + " VARCHAR(88) NOT NULL, " + H_struct.VECTOR_CLOCK + " VARCHAR(88) NOT NULL, " +
@ -193,22 +212,37 @@ function Database(user, password, dbname, host = 'localhost') {
F_struct.NOTE_SIGN + " VARCHAR(160), " + F_struct.NOTE_SIGN + " VARCHAR(160), " +
"PRIMARY KEY (" + H_struct.VECTOR_CLOCK + ")" + "PRIMARY KEY (" + H_struct.VECTOR_CLOCK + ")" +
" )"; " )";
db.query(statement) queryResolve(statement)
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.dropTable = function (snID) { DB.dropTable = function (snID) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "DROP TABLE _" + snID; let statement = "DROP TABLE _" + snID;
db.query(statement) queryResolve(statement)
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.addData = function (snID, data) { DB.listTable = function () {
return new Promise((resolve, reject) => {
queryResolve("SHOW TABLES").then(result => {
const disks = [];
for (let i in result)
for (let j in result[i])
if (result[i][j].startsWith("_"))
disks.push(result[i][j].split("_")[1]);
resolve(disks);
}).catch(error => reject(error))
})
}
/* Data Service (by client)*/
DB.addData = function (snID, data) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
data[L_struct.STATUS] = 1; data[L_struct.STATUS] = 1;
data[L_struct.LOG_TIME] = Date.now(); data[L_struct.LOG_TIME] = Date.now();
@ -222,23 +256,23 @@ function Database(user, password, dbname, host = 'localhost') {
" (" + attr.join(", ") + ") " + " (" + attr.join(", ") + ") " +
"VALUES (" + attr.map(a => '?').join(", ") + ")"; "VALUES (" + attr.map(a => '?').join(", ") + ")";
data = Object.fromEntries(attr.map((a, i) => [a, values[i]])); data = Object.fromEntries(attr.map((a, i) => [a, values[i]]));
db.query(statement, values) queryResolve(statement, values)
.then(result => resolve(data)) .then(result => resolve(data))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.getData = function (snID, vectorClock) { DB.getData = function (snID, vectorClock) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "SELECT * FROM _" + snID + let statement = "SELECT * FROM _" + snID +
" WHERE " + H_struct.VECTOR_CLOCK + "=?"; " WHERE " + H_struct.VECTOR_CLOCK + "=?";
db.query(statement, [vectorClock]) queryResolve(statement, [vectorClock])
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)) .catch(error => reject(error))
}) })
}; };
db.tagData = function (snID, vectorClock, tag, tagTime, tagKey, tagSign) { DB.tagData = function (snID, vectorClock, tag, tagTime, tagKey, tagSign) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let data = { let data = {
[T_struct.TAG]: tag, [T_struct.TAG]: tag,
@ -253,13 +287,13 @@ function Database(user, password, dbname, host = 'localhost') {
let statement = "UPDATE _" + snID + let statement = "UPDATE _" + snID +
" SET " + attr.map(a => a + "=?").join(", ") + " SET " + attr.map(a => a + "=?").join(", ") +
" WHERE " + H_struct.VECTOR_CLOCK + "=?"; " WHERE " + H_struct.VECTOR_CLOCK + "=?";
db.query(statement, values) queryResolve(statement, values)
.then(result => resolve(data)) .then(result => resolve(data))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.noteData = function (snID, vectorClock, note, noteTime, noteKey, noteSign) { DB.noteData = function (snID, vectorClock, note, noteTime, noteKey, noteSign) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let data = { let data = {
[F_struct.NOTE]: note, [F_struct.NOTE]: note,
@ -274,13 +308,13 @@ function Database(user, password, dbname, host = 'localhost') {
let statement = "UPDATE _" + snID + let statement = "UPDATE _" + snID +
" SET " + attr.map(a => a + "=?").join(", ") + " SET " + attr.map(a => a + "=?").join(", ") +
" WHERE " + H_struct.VECTOR_CLOCK + "=?"; " WHERE " + H_struct.VECTOR_CLOCK + "=?";
db.query(statement, values) queryResolve(statement, values)
.then(result => resolve(data)) .then(result => resolve(data))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
} }
db.searchData = function (snID, request) { DB.searchData = function (snID, request) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let conditionArr = []; let conditionArr = [];
if (request.lowerVectorClock || request.upperVectorClock || request.atVectorClock) { if (request.lowerVectorClock || request.upperVectorClock || request.atVectorClock) {
@ -316,44 +350,46 @@ function Database(user, password, dbname, host = 'localhost') {
" WHERE " + conditionArr.join(" AND ") + " WHERE " + conditionArr.join(" AND ") +
" ORDER BY " + (request.afterTime ? L_struct.LOG_TIME : H_struct.VECTOR_CLOCK) + " ORDER BY " + (request.afterTime ? L_struct.LOG_TIME : H_struct.VECTOR_CLOCK) +
(request.mostRecent ? " DESC LIMIT 1" : ""); (request.mostRecent ? " DESC LIMIT 1" : "");
db.query(statement) queryResolve(statement)
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.lastLogTime = function (snID) { /* Backup service */
DB.lastLogTime = function (snID) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let attr = "MAX(" + L_struct.LOG_TIME + ")"; let attr = "MAX(" + L_struct.LOG_TIME + ")";
let statement = "SELECT " + attr + " FROM _" + snID; let statement = "SELECT " + attr + " FROM _" + snID;
db.query(statement) queryResolve(statement)
.then(result => resolve(result[0][attr] || 0)) .then(result => resolve(result[0][attr] || 0))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.createGetLastLog = function (snID) { DB.createGetLastLog = function (snID) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
db.createTable(snID).then(result => { DB.createTable(snID).then(result => {
db.lastLogTime(snID) DB.lastLogTime(snID)
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}).catch(error => reject(error)); }).catch(error => reject(error));
}); });
}; };
db.readAllDataStream = function (snID, logtime, callback) { DB.readAllDataStream = function (snID, logtime, callback) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "SELECT * FROM _" + snID + let statement = "SELECT * FROM _" + snID +
" WHERE " + L_struct.LOG_TIME + ">" + logtime + " WHERE " + L_struct.LOG_TIME + ">" + logtime +
" ORDER BY " + L_struct.LOG_TIME; " ORDER BY " + L_struct.LOG_TIME;
db.query_stream(statement, callback) queryStream(statement, callback)
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.storeData = function (snID, data, updateLogTime = false) { DB.storeData = function (snID, data, updateLogTime = false) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
if (updateLogTime) if (updateLogTime)
data[L_struct.LOG_TIME] = Date.now(); data[L_struct.LOG_TIME] = Date.now();
@ -367,47 +403,49 @@ function Database(user, password, dbname, host = 'localhost') {
" (" + attr.join(", ") + ") " + " (" + attr.join(", ") + ") " +
"VALUES (" + attr.map(a => '?').join(", ") + ") " + "VALUES (" + attr.map(a => '?').join(", ") + ") " +
"ON DUPLICATE KEY UPDATE " + u_attr.map(a => a + "=?").join(", "); "ON DUPLICATE KEY UPDATE " + u_attr.map(a => a + "=?").join(", ");
db.query(statement, values.concat(u_values)) queryResolve(statement, values.concat(u_values))
.then(result => resolve(data)) .then(result => resolve(data))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.storeTag = function (snID, data) { DB.storeTag = function (snID, data) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let attr = Object.keys(T_struct).map(a => T_struct[a]).concat(L_struct.LOG_TIME); let attr = Object.keys(T_struct).map(a => T_struct[a]).concat(L_struct.LOG_TIME);
let values = attr.map(a => data[a]).concat(data[H_struct.VECTOR_CLOCK]); let values = attr.map(a => data[a]).concat(data[H_struct.VECTOR_CLOCK]);
let statement = "UPDATE _" + snID + let statement = "UPDATE _" + snID +
" SET " + attr.map(a => a + "=?").join(", ") + " SET " + attr.map(a => a + "=?").join(", ") +
" WHERE " + H_struct.VECTOR_CLOCK + "=?"; " WHERE " + H_struct.VECTOR_CLOCK + "=?";
db.query(statement, values) queryResolve(statement, values)
.then(result => resolve(data)) .then(result => resolve(data))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.storeNote = function (snID, data) { DB.storeNote = function (snID, data) {
let attr = Object.keys(F_struct).map(a => F_struct[a]).concat(L_struct.LOG_TIME); let attr = Object.keys(F_struct).map(a => F_struct[a]).concat(L_struct.LOG_TIME);
let values = attr.map(a => data[a]).concat(data[H_struct.VECTOR_CLOCK]); let values = attr.map(a => data[a]).concat(data[H_struct.VECTOR_CLOCK]);
let statement = "UPDATE _" + snID + let statement = "UPDATE _" + snID +
" SET " + attr.map(a => a + "=?").join(", ") + " SET " + attr.map(a => a + "=?").join(", ") +
" WHERE " + H_struct.VECTOR_CLOCK + "=?"; " WHERE " + H_struct.VECTOR_CLOCK + "=?";
db.query(statement, values) queryResolve(statement, values)
.then(result => resolve(data)) .then(result => resolve(data))
.catch(error => reject(error)); .catch(error => reject(error));
} }
db.deleteData = function (snID, vectorClock) { DB.deleteData = function (snID, vectorClock) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "DELETE FROM _" + snID + let statement = "DELETE FROM _" + snID +
" WHERE " + H_struct.VECTOR_CLOCK + "=?"; " WHERE " + H_struct.VECTOR_CLOCK + "=?";
db.query(statement, [vectorClock]) queryResolve(statement, [vectorClock])
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.clearAuthorisedAppData = function (snID, app, adminID, subAdmins, timestamp) { /* Data clearing */
DB.clearAuthorisedAppData = function (snID, app, adminID, subAdmins, timestamp) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "DELETE FROM _" + snID + let statement = "DELETE FROM _" + snID +
" WHERE ( " + H_struct.TIME + "<? AND " + " WHERE ( " + H_struct.TIME + "<? AND " +
@ -417,86 +455,28 @@ function Database(user, password, dbname, host = 'localhost') {
H_struct.RECEIVER_ID + " != ? OR " + H_struct.RECEIVER_ID + " != ? OR " +
H_struct.SENDER_ID + " NOT IN (" + subAdmins.map(a => "?").join(", ") + ") )" : H_struct.SENDER_ID + " NOT IN (" + subAdmins.map(a => "?").join(", ") + ") )" :
""); "");
db.query(statement, [timestamp, app, adminID].concat(subAdmins)) queryResolve(statement, [timestamp, app, adminID].concat(subAdmins))
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
}; };
db.clearUnauthorisedAppData = function (snID, appList, timestamp) { DB.clearUnauthorisedAppData = function (snID, appList, timestamp) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let statement = "DELETE FROM _" + snID + let statement = "DELETE FROM _" + snID +
" WHERE " + H_struct.TIME + "<?" + " WHERE " + H_struct.TIME + "<?" +
(appList.length ? " AND " + (appList.length ? " AND " +
H_struct.APPLICATION + " NOT IN (" + appList.map(a => "?").join(", ") + ")" : H_struct.APPLICATION + " NOT IN (" + appList.map(a => "?").join(", ") + ")" :
""); "");
db.query(statement, [timestamp].concat(appList)) queryResolve(statement, [timestamp].concat(appList))
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}); });
};
Object.defineProperty(db, "connect", {
get: () => new Promise((resolve, reject) => {
db.pool.getConnection((error, conn) => {
if (error)
reject(error);
else
resolve(conn);
});
})
});
Object.defineProperty(db, "query", {
value: (sql, values) => new Promise((resolve, reject) => {
db.connect.then(conn => {
const fn = (err, res) => {
conn.release();
(err ? reject(err) : resolve(res));
};
if (values)
conn.query(sql, values, fn);
else
conn.query(sql, fn);
}).catch(error => reject(error));
})
});
Object.defineProperty(db, "query_stream", {
value: (sql, value, callback) => new Promise((resolve, reject) => {
db.connect.then(conn => {
if (!callback && value instanceof Function) {
callback = value;
value = undefined;
}
var err_flag, row_count = 0;
(value ? conn.query(sql, values) : conn.query(sql))
.on('error', err => err_flag = err)
.on('result', row => {
row_count++;
conn.pause();
callback(row);
conn.resume();
}).on('end', _ => {
conn.release();
err_flag ? reject(err_flag) : resolve(row_count);
});
})
})
})
return new Promise((resolve, reject) => {
db.pool = mysql.createPool({
host: host,
user: user,
password: password,
database: dbname
});
db.connect.then(conn => {
conn.release();
resolve(db);
}).catch(error => reject(error));
});
}; };
module.exports = Database; module.exports = {
init: initConnection,
query: queryResolve,
query_stream: queryStream,
DB
};

View File

@ -6,11 +6,10 @@ global.floCrypto = require('./floCrypto');
global.floBlockchainAPI = require('./floBlockchainAPI'); global.floBlockchainAPI = require('./floBlockchainAPI');
const Database = require("./database"); const Database = require("./database");
const intra = require('./intra'); const intra = require('./intra');
const client = require('./client');
const Server = require('./server'); const Server = require('./server');
const keys = require("./keys"); const keys = require("./keys");
var DB; //Container for Database object const DB = Database.DB;
const INTERVAL_REFRESH_TIME = 1 * 60 * 60 * 1000; //1 hr const INTERVAL_REFRESH_TIME = 1 * 60 * 60 * 1000; //1 hr
function startNode() { function startNode() {
@ -36,13 +35,8 @@ function startNode() {
console.info("Logged in as", keys.node_id); console.info("Logged in as", keys.node_id);
//DB connect //DB connect
Database(config["sql_user"], config["sql_pwd"], config["sql_db"], config["sql_host"]).then(db => { Database.init(config["sql_user"], config["sql_pwd"], config["sql_db"], config["sql_host"]).then(db => {
console.info("Connected to Database"); console.info("Connected to Database");
DB = db;
//Set DB to client and intra scripts
intra.DB = DB;
client.DB = DB;
client._list = intra._list;
loadBase().then(base => { loadBase().then(base => {
console.log("Load Database successful"); console.log("Load Database successful");
//Set base data from DB to floGlobals //Set base data from DB to floGlobals
@ -54,7 +48,7 @@ function startNode() {
refreshData.invoke(null) refreshData.invoke(null)
.then(_ => intra.reconnectNextNode()).catch(_ => null); .then(_ => intra.reconnectNextNode()).catch(_ => null);
//Start Server //Start Server
const server = new Server(config["port"], client, intra); const server = new Server(config["port"]);
server.refresher = refreshData; server.refresher = refreshData;
intra.refresher = refreshData; intra.refresher = refreshData;
}).catch(error => console.error(error)); }).catch(error => console.error(error));
@ -64,7 +58,7 @@ function startNode() {
function loadBase() { function loadBase() {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
DB.createBase().then(result => { DB.createBase().then(result => {
DB.getBase(DB) DB.getBase()
.then(result => resolve(result)) .then(result => resolve(result))
.catch(error => reject(error)); .catch(error => reject(error));
}).catch(error => reject(error)); }).catch(error => reject(error));
@ -257,12 +251,7 @@ function diskCleanUp(base) {
}; };
function selfDiskMigration(node_change) { function selfDiskMigration(node_change) {
DB.query("SHOW TABLES").then(result => { DB.listTable("SHOW TABLES").then(disks => {
const disks = [];
for (let i in result)
for (let j in result[i])
if (result[i][j].startsWith("_"))
disks.push(result[i][j].split("_")[1]);
disks.forEach(n => { disks.forEach(n => {
if (node_change[n] === false) if (node_change[n] === false)
DB.dropTable(n).then(_ => null).catch(e => console.error(e)); DB.dropTable(n).then(_ => null).catch(e => console.error(e));

View File

@ -1,11 +1,13 @@
const http = require('http'); const http = require('http');
const WebSocket = require('ws'); const WebSocket = require('ws');
const url = require('url'); const url = require('url');
const intra = require('./intra');
const client = require('./client');
const INVALID_E_CODE = 400, const INVALID_E_CODE = 400,
INTERNAL_E_CODE = 500; INTERNAL_E_CODE = 500;
module.exports = function Server(port, client, intra) { module.exports = function Server(port) {
var refresher; //container for refresher var refresher; //container for refresher
@ -62,7 +64,7 @@ module.exports = function Server(port, client, intra) {
server server
}); });
wsServer.on('connection', function connection(ws) { wsServer.on('connection', function connection(ws) {
ws.onmessage = function(evt) { ws.onmessage = function (evt) {
let message = evt.data; let message = evt.data;
if (message.startsWith(intra.SUPERNODE_INDICATOR)) if (message.startsWith(intra.SUPERNODE_INDICATOR))
intra.processTaskFromSupernode(message, ws); intra.processTaskFromSupernode(message, ws);