mirror of
https://github.com/cryptpad/cryptpad.git
synced 2026-09-12 19:49:59 +05:00
Add aggregated stats for instances that opted-in
This commit is contained in:
parent
9d5af0d136
commit
fd774282c9
61
lib/api.js
61
lib/api.js
@ -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);
|
||||||
|
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
@ -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});
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
};
|
};
|
||||||
|
|||||||
@ -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 ||
|
||||||
|
|||||||
@ -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;
|
||||||
|
|||||||
@ -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();
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user