mirror of
https://github.com/cryptpad/cryptpad.git
synced 2026-09-12 19:49:59 +05:00
Merge pull request #1509 from cryptpad/stats
Server fixes and aggregated stats
This commit is contained in:
commit
3a0fd0ecd6
61
lib/api.js
61
lib/api.js
@ -61,7 +61,7 @@ nThen(function (w) {
|
||||
};
|
||||
|
||||
// spawn ws server and attach netflux event handlers
|
||||
NetfluxSrv.create(new WebSocketServer({ server: Env.httpServer}))
|
||||
let Server = NetfluxSrv.create(new WebSocketServer({ server: Env.httpServer}))
|
||||
.on('channelClose', historyKeeper.channelClose)
|
||||
.on('channelMessage', historyKeeper.channelMessage)
|
||||
.on('channelOpen', historyKeeper.channelOpen)
|
||||
@ -90,6 +90,65 @@ nThen(function (w) {
|
||||
});
|
||||
})
|
||||
.register(historyKeeper.id, historyKeeper.directMessage);
|
||||
// Store max active WS during the last day (reset when sending ping if enabled)
|
||||
setInterval(() => {
|
||||
try {
|
||||
// Concurrent usage data
|
||||
let oldWs = Env.maxConcurrentWs || 0;
|
||||
let oldUniqueWs = Env.maxConcurrentUniqueWs || 0;
|
||||
let oldChans = Env.maxActiveChannels || 0;
|
||||
let oldUsers = Env.maxConcurrentRegUsers || 0;
|
||||
let stats = Server.getSessionStats();
|
||||
let chans = Server.getActiveChannelCount();
|
||||
let reg = 0;
|
||||
let regKeys = [];
|
||||
Object.keys(Env.netfluxUsers).forEach(id => {
|
||||
let keys = Env.netfluxUsers[id];
|
||||
let key = Object.keys(keys || {})[0];
|
||||
if (!key) { return; }
|
||||
if (regKeys.includes(key)) { return; }
|
||||
reg++;
|
||||
regKeys.push(key);
|
||||
});
|
||||
Env.maxConcurrentWs = Math.max(oldWs, stats.total);
|
||||
Env.maxConcurrentUniqueWs = Math.max(oldUniqueWs, stats.unique);
|
||||
Env.maxConcurrentRegUsers = Math.max(oldUsers, reg);
|
||||
Env.maxActiveChannels = Math.max(oldChans, chans);
|
||||
} catch (e) {}
|
||||
}, 10000);
|
||||
// Clean up active registered users and channels (possible memory leak)
|
||||
setInterval(() => {
|
||||
try {
|
||||
let users = Env.netfluxUsers || {};
|
||||
let online = Server.getOnlineUsers() || [];
|
||||
let removed = 0;
|
||||
Object.keys(users).forEach(id => {
|
||||
if (!online.includes(id)) {
|
||||
delete users[id];
|
||||
removed++;
|
||||
}
|
||||
});
|
||||
if (removed) {
|
||||
Env.Log.info("CLEANED_ACTIVE_USERS_MAP", {removed});
|
||||
}
|
||||
} catch (e) {}
|
||||
try {
|
||||
let HK = require('./hk-utils');
|
||||
let chans = Env.channel_cache || {};
|
||||
let active = Server.getActiveChannels() || [];
|
||||
let removed = 0;
|
||||
Object.keys(chans).forEach(id => {
|
||||
if (!active.includes(id)) {
|
||||
HK.dropChannel(Env, id);
|
||||
removed++;
|
||||
}
|
||||
});
|
||||
if (removed) {
|
||||
Env.Log.info("CLEANED_ACTIVE_CHANNELS_MAP", {removed});
|
||||
}
|
||||
} catch (e) {}
|
||||
}, 30000);
|
||||
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
@ -102,11 +102,13 @@ var shutdown = function (Env, Server, cb) {
|
||||
// and allow system functionality to restart the server
|
||||
};
|
||||
|
||||
var getRegisteredUsers = function (Env, Server, cb) {
|
||||
var getRegisteredUsers = Admin.getRegisteredUsers = function (Env, Server, cb) {
|
||||
Env.batchRegisteredUsers('', cb, function (done) {
|
||||
var dir = Env.paths.pin;
|
||||
var folders;
|
||||
var dirB = Env.paths.block;
|
||||
var folders, foldersB;
|
||||
var users = 0;
|
||||
var blocks = 0;
|
||||
nThen(function (waitFor) {
|
||||
Fs.readdir(dir, waitFor(function (err, list) {
|
||||
if (err) {
|
||||
@ -115,16 +117,39 @@ var getRegisteredUsers = function (Env, Server, cb) {
|
||||
}
|
||||
folders = list;
|
||||
}));
|
||||
Fs.readdir(dirB, waitFor(function (err, list) {
|
||||
if (err) {
|
||||
waitFor.abort();
|
||||
return void done(err);
|
||||
}
|
||||
foldersB = list;
|
||||
}));
|
||||
}).nThen(function (waitFor) {
|
||||
folders.forEach(function (f) {
|
||||
var dir = Env.paths.pin + '/' + f;
|
||||
Fs.readdir(dir, waitFor(function (err, list) {
|
||||
if (err) { return; }
|
||||
// Don't count placeholders
|
||||
list = list.filter(name => {
|
||||
return !/\.placeholder$/.test(name);
|
||||
});
|
||||
users += list.length;
|
||||
}));
|
||||
});
|
||||
}).nThen(function (waitFor) {
|
||||
foldersB.forEach(function (f) {
|
||||
var dir = Env.paths.block + '/' + f;
|
||||
Fs.readdir(dir, waitFor(function (err, list) {
|
||||
if (err) { return; }
|
||||
// Don't count placeholders
|
||||
list = list.filter(name => {
|
||||
return !/\.placeholder$/.test(name);
|
||||
});
|
||||
blocks += list.length;
|
||||
}));
|
||||
});
|
||||
}).nThen(function () {
|
||||
done(void 0, users);
|
||||
done(void 0, {users, blocks});
|
||||
});
|
||||
});
|
||||
};
|
||||
|
||||
@ -11,6 +11,8 @@ const Https = require("https");
|
||||
const Http = require("http");
|
||||
const Util = require("../common-util");
|
||||
const Stats = require("../stats");
|
||||
const Admin = require("./admin-rpc.js");
|
||||
const nThen = require('nthen');
|
||||
|
||||
var validLimitFields = ['limit', 'plan', 'note', 'users', 'origin'];
|
||||
|
||||
@ -112,45 +114,85 @@ var queryAccountServer = function (Env, cb) {
|
||||
var done = Util.once(Util.mkAsync(cb));
|
||||
|
||||
var rawBody = Stats.instanceData(Env);
|
||||
Env.Log.info("SERVER_TELEMETRY", rawBody);
|
||||
var body = JSON.stringify(rawBody);
|
||||
|
||||
var options = {
|
||||
host: 'accounts.cryptpad.fr',
|
||||
path: '/api/getauthorized',
|
||||
method: 'POST',
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
"Content-Length": Buffer.byteLength(body)
|
||||
}
|
||||
let send = () => {
|
||||
Env.Log.info("SERVER_TELEMETRY", rawBody);
|
||||
var body = JSON.stringify(rawBody);
|
||||
|
||||
var options = {
|
||||
host: 'accounts.cryptpad.fr',
|
||||
path: '/api/getauthorized',
|
||||
method: 'POST',
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
"Content-Length": Buffer.byteLength(body)
|
||||
}
|
||||
};
|
||||
|
||||
var req = Https.request(options, function (response) {
|
||||
if (!('' + response.statusCode).match(/^2\d\d$/)) {
|
||||
return void cb('SERVER ERROR ' + response.statusCode);
|
||||
}
|
||||
var str = '';
|
||||
|
||||
response.on('data', function (chunk) {
|
||||
str += chunk;
|
||||
});
|
||||
|
||||
response.on('end', function () {
|
||||
try {
|
||||
var json = JSON.parse(str);
|
||||
checkUpdateAvailability(Env, json);
|
||||
done(void 0, json);
|
||||
} catch (e) {
|
||||
done(e);
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
req.on('error', function () {
|
||||
done();
|
||||
});
|
||||
|
||||
req.end(body);
|
||||
};
|
||||
|
||||
var req = Https.request(options, function (response) {
|
||||
if (!('' + response.statusCode).match(/^2\d\d$/)) {
|
||||
return void cb('SERVER ERROR ' + response.statusCode);
|
||||
}
|
||||
var str = '';
|
||||
|
||||
response.on('data', function (chunk) {
|
||||
str += chunk;
|
||||
});
|
||||
|
||||
response.on('end', function () {
|
||||
try {
|
||||
var json = JSON.parse(str);
|
||||
checkUpdateAvailability(Env, json);
|
||||
done(void 0, json);
|
||||
} catch (e) {
|
||||
done(e);
|
||||
if (Env.provideAggregateStatistics) {
|
||||
let stats = {};
|
||||
nThen(waitFor => {
|
||||
Admin.getRegisteredUsers(Env, null, waitFor((err, data) => {
|
||||
if (err) { return; }
|
||||
stats.registered = data.blocks;
|
||||
if (Env.lastPingRegisteredUsers) {
|
||||
stats.usersDiff = stats.registered - Env.lastPingRegisteredUsers;
|
||||
}
|
||||
Env.lastPingRegisteredUsers = stats.registered;
|
||||
let teams = (data.users - data.blocks);
|
||||
if (teams > 0) { stats.teams = teams; }
|
||||
}));
|
||||
}).nThen(() => {
|
||||
if (Env.maxConcurrentWs) {
|
||||
stats.maxConcurrentWs = Env.maxConcurrentWs;
|
||||
Env.maxConcurrentWs = 0;
|
||||
}
|
||||
if (Env.maxConcurrentUniqueWs) {
|
||||
stats.maxConcurrentUniqueIPs = Env.maxConcurrentUniqueWs;
|
||||
Env.maxConcurrentUniqueWs = 0;
|
||||
}
|
||||
if (Env.maxConcurrentRegUsers) {
|
||||
stats.maxConcurrentRegUsers = Env.maxConcurrentRegUsers;
|
||||
Env.maxConcurrentRegUsers = 0;
|
||||
}
|
||||
if (Env.maxActiveChannels) {
|
||||
stats.maxConcurrentChannels = Env.maxActiveChannels;
|
||||
Env.maxActiveChannels = 0;
|
||||
}
|
||||
rawBody.statistics = stats;
|
||||
send();
|
||||
});
|
||||
});
|
||||
|
||||
req.on('error', function () {
|
||||
done();
|
||||
});
|
||||
|
||||
req.end(body);
|
||||
return;
|
||||
}
|
||||
send();
|
||||
};
|
||||
Quota.shouldContactServer = function (Env) {
|
||||
return !(Env.blockDailyCheck === true ||
|
||||
|
||||
@ -81,6 +81,7 @@ Stats.instanceData = function (Env) {
|
||||
if (Env.provideAggregateStatistics) {
|
||||
// check how many instances provide stats before we put more work into it
|
||||
data.providesAggregateStatistics = true;
|
||||
data.statistics = {}; // Filled in lib/commands/quota.js because of async calls
|
||||
}
|
||||
|
||||
return data;
|
||||
|
||||
@ -2647,9 +2647,11 @@ define([
|
||||
var onRefresh = function () {
|
||||
sFrameChan.query('Q_ADMIN_RPC', {
|
||||
cmd: 'REGISTERED_USERS',
|
||||
}, function (e, data) {
|
||||
}, function (e, arr) {
|
||||
pre.innerText = '';
|
||||
pre.append(String(data));
|
||||
let data = arr[0];
|
||||
pre.append(String(data.blocks));
|
||||
pre.append(' (old value including teams: ' + String(data.users) + ')'); // XXX
|
||||
});
|
||||
};
|
||||
onRefresh();
|
||||
|
||||
Loading…
Reference in New Issue
Block a user