Add aggregated stats for instances that opted-in

This commit is contained in:
yflory 2024-06-11 18:14:45 +02:00
parent 9d5af0d136
commit fd774282c9
5 changed files with 169 additions and 40 deletions

View File

@ -61,7 +61,7 @@ nThen(function (w) {
}; };
// spawn ws server and attach netflux event handlers // 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('channelClose', historyKeeper.channelClose)
.on('channelMessage', historyKeeper.channelMessage) .on('channelMessage', historyKeeper.channelMessage)
.on('channelOpen', historyKeeper.channelOpen) .on('channelOpen', historyKeeper.channelOpen)
@ -90,6 +90,65 @@ nThen(function (w) {
}); });
}) })
.register(historyKeeper.id, historyKeeper.directMessage); .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);
}); });
}); });

View File

@ -102,11 +102,13 @@ var shutdown = function (Env, Server, cb) {
// and allow system functionality to restart the server // 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) { Env.batchRegisteredUsers('', cb, function (done) {
var dir = Env.paths.pin; var dir = Env.paths.pin;
var folders; var dirB = Env.paths.block;
var folders, foldersB;
var users = 0; var users = 0;
var blocks = 0;
nThen(function (waitFor) { nThen(function (waitFor) {
Fs.readdir(dir, waitFor(function (err, list) { Fs.readdir(dir, waitFor(function (err, list) {
if (err) { if (err) {
@ -115,16 +117,39 @@ var getRegisteredUsers = function (Env, Server, cb) {
} }
folders = list; folders = list;
})); }));
Fs.readdir(dirB, waitFor(function (err, list) {
if (err) {
waitFor.abort();
return void done(err);
}
foldersB = list;
}));
}).nThen(function (waitFor) { }).nThen(function (waitFor) {
folders.forEach(function (f) { folders.forEach(function (f) {
var dir = Env.paths.pin + '/' + f; var dir = Env.paths.pin + '/' + f;
Fs.readdir(dir, waitFor(function (err, list) { Fs.readdir(dir, waitFor(function (err, list) {
if (err) { return; } if (err) { return; }
// Don't count placeholders
list = list.filter(name => {
return !/\.placeholder$/.test(name);
});
users += list.length; 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 () { }).nThen(function () {
done(void 0, users); done(void 0, {users, blocks});
}); });
}); });
}; };

View File

@ -11,6 +11,8 @@ const Https = require("https");
const Http = require("http"); const Http = require("http");
const Util = require("../common-util"); const Util = require("../common-util");
const Stats = require("../stats"); const Stats = require("../stats");
const Admin = require("./admin-rpc.js");
const nThen = require('nthen');
var validLimitFields = ['limit', 'plan', 'note', 'users', 'origin']; var validLimitFields = ['limit', 'plan', 'note', 'users', 'origin'];
@ -112,45 +114,85 @@ var queryAccountServer = function (Env, cb) {
var done = Util.once(Util.mkAsync(cb)); var done = Util.once(Util.mkAsync(cb));
var rawBody = Stats.instanceData(Env); var rawBody = Stats.instanceData(Env);
Env.Log.info("SERVER_TELEMETRY", rawBody);
var body = JSON.stringify(rawBody);
var options = { let send = () => {
host: 'accounts.cryptpad.fr', Env.Log.info("SERVER_TELEMETRY", rawBody);
path: '/api/getauthorized', var body = JSON.stringify(rawBody);
method: 'POST',
headers: { var options = {
"Content-Type": "application/json", host: 'accounts.cryptpad.fr',
"Content-Length": Buffer.byteLength(body) 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 (Env.provideAggregateStatistics) {
if (!('' + response.statusCode).match(/^2\d\d$/)) { let stats = {};
return void cb('SERVER ERROR ' + response.statusCode); nThen(waitFor => {
} Admin.getRegisteredUsers(Env, null, waitFor((err, data) => {
var str = ''; if (err) { return; }
stats.registered = data.blocks;
response.on('data', function (chunk) { if (Env.lastPingRegisteredUsers) {
str += chunk; stats.usersDiff = stats.registered - Env.lastPingRegisteredUsers;
}); }
Env.lastPingRegisteredUsers = stats.registered;
response.on('end', function () { let teams = (data.users - data.blocks);
try { if (teams > 0) { stats.teams = teams; }
var json = JSON.parse(str); }));
checkUpdateAvailability(Env, json); }).nThen(() => {
done(void 0, json); if (Env.maxConcurrentWs) {
} catch (e) { stats.maxConcurrentWs = Env.maxConcurrentWs;
done(e); 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();
}); });
}); return;
}
req.on('error', function () { send();
done();
});
req.end(body);
}; };
Quota.shouldContactServer = function (Env) { Quota.shouldContactServer = function (Env) {
return !(Env.blockDailyCheck === true || return !(Env.blockDailyCheck === true ||

View File

@ -81,6 +81,7 @@ Stats.instanceData = function (Env) {
if (Env.provideAggregateStatistics) { if (Env.provideAggregateStatistics) {
// check how many instances provide stats before we put more work into it // check how many instances provide stats before we put more work into it
data.providesAggregateStatistics = true; data.providesAggregateStatistics = true;
data.statistics = {}; // Filled in lib/commands/quota.js because of async calls
} }
return data; return data;

View File

@ -2647,9 +2647,11 @@ define([
var onRefresh = function () { var onRefresh = function () {
sFrameChan.query('Q_ADMIN_RPC', { sFrameChan.query('Q_ADMIN_RPC', {
cmd: 'REGISTERED_USERS', cmd: 'REGISTERED_USERS',
}, function (e, data) { }, function (e, arr) {
pre.innerText = ''; 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(); onRefresh();