diff --git a/lib/api.js b/lib/api.js index 220e954d8..91ae1598b 100644 --- a/lib/api.js +++ b/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); + }); }); diff --git a/lib/commands/admin-rpc.js b/lib/commands/admin-rpc.js index 954c1c4c1..6661bd25b 100644 --- a/lib/commands/admin-rpc.js +++ b/lib/commands/admin-rpc.js @@ -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}); }); }); }; diff --git a/lib/commands/quota.js b/lib/commands/quota.js index 7a92dbda9..e53779b5d 100644 --- a/lib/commands/quota.js +++ b/lib/commands/quota.js @@ -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 || diff --git a/lib/stats.js b/lib/stats.js index d67a2ee71..d8d1d00ba 100644 --- a/lib/stats.js +++ b/lib/stats.js @@ -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; diff --git a/www/admin/inner.js b/www/admin/inner.js index cd1360d9c..bb930e247 100644 --- a/www/admin/inner.js +++ b/www/admin/inner.js @@ -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();