fix: replace shared telnet singleton with connection pool

Addresses PERFORMANCE-FINDINGS #1, #4, #6 (and #2 in creategroup flow):
- Add app/telnet-pool.js: pool of N connections (config.telnet.poolSize,
  default 4) with wait-queue, lazy reconnect, and respawn on dead socket.
- Remove module-level new Telnet() singleton and startup socket dump from
  app/routes/index.js; switch all mutating endpoints to telnetPool.send().
- Set telnet debug:false; use config.telnet.timeout (default 5s) instead
  of hardcoded 30000.
- Add telnet.poolSize to config/prosody-muc-rest.js.
- await doesMucExist in /creategroup so the existence check is no longer dead.

Part-of: <http://gitlab.vnc.biz/uxf/prosody-muc-rest/-/merge_requests/3>
This commit is contained in:
2026-07-10 14:44:43 +02:00
parent 646f0088cd
commit 1fad2ef1d2
3 changed files with 195 additions and 146 deletions
+94 -145
View File
@@ -1,30 +1,21 @@
var express = require("express");
var router = express.Router();
var Telnet = require('telnet-client');
var moment = require('moment');
var request = require('@cypress/request');
var env = process.env.NODE_ENV || 'development';
var config = require('../../config/prosody-muc-rest.js')[env];
var Pool = require('pg-pool');
var createTelnetPool = require('../telnet-pool');
var connection = new Telnet();
let params = {
var telnetParams = {
host: config.telnet.host,
port: config.telnet.port,
// shellPrompt: '', // or negotiationMandatory: false
negotiationMandatory: config.telnet.negotiationMandatory,
// timeout: config.telnet.timeout,
timeout: 30000,
debug: true
timeout: config.telnet.timeout,
debug: false
}
try {
connection.connect(params);
console.log("connection is now: ", connection);
} catch(error) {
// handle the throw (timeout)
console.log("error while connect to telnet:", error);
}
var telnetPool = createTelnetPool(telnetParams, config.telnet.poolSize);
async function doesMucExist(jid) {
return new Promise((resolve, reject) => {
@@ -282,32 +273,22 @@ router.put('/groupchats/:target', async function (req, res) {
}
if (has_data || has_subject) {
if (has_subject) {
if ( connection && connection.socket) {
if (!connection.socket.writable) {
console.log("connection not writable - trying to reconnect");
await connection.connect(params).catch(error => {
console.log("connection error: ", error);
});
console.log("reconnect result: ", connection.socket.writable);
}
if ( connection.socket.writable === false ) {
res.status(500).json({message: 'telnet connection is not established.'});
} else {
let response = await connection.send('muc:room(\"'+ target +'\"):set_subject(\"'+ user +'\", \"'+ subject +'\");\n');
let rawresults = response.split("\n");
let mucresults = [];
for (var i = 0; i < rawresults.length - 1 ; i++) {
if ((rawresults[i].startsWith("\u0000| ") || rawresults[i].startsWith("\| ")) && !rawresults[i].startsWith("| OK") && !rawresults[i].startsWith("\u0000| Result") ) {
console.log("adding: ", rawresults[i].split("| ")[1].split("\r")[0]);
mucresults.push(rawresults[i].split("| ")[1].split("\r")[0]);
}
try {
let response = await telnetPool.send('muc:room(\"'+ target +'\"):set_subject(\"'+ user +'\", \"'+ subject +'\");\n');
let rawresults = response.split("\n");
let mucresults = [];
for (var i = 0; i < rawresults.length - 1 ; i++) {
if ((rawresults[i].startsWith("\u0000| ") || rawresults[i].startsWith("\| ")) && !rawresults[i].startsWith("| OK") && !rawresults[i].startsWith("\u0000| Result") ) {
console.log("adding: ", rawresults[i].split("| ")[1].split("\r")[0]);
mucresults.push(rawresults[i].split("| ")[1].split("\r")[0]);
}
console.log('async rawresult:', rawresults);
console.log('async result:', mucresults);
res.status(200).json(mucresults);
}
} else {
res.status(404).json({message: 'connection not available'});
console.log('async rawresult:', rawresults);
console.log('async result:', mucresults);
res.status(200).json(mucresults);
} catch (error) {
console.log("telnet send error: ", error);
res.status(500).json({message: 'telnet connection is not established.'});
}
} else {
if (has_data) {
@@ -434,37 +415,26 @@ router.post('/affiliations/:target', async function (req, res) {
if (derr == null) {
if (dres.rowCount > 0) {
// groupchat is found and exists
if ( connection && connection.socket) {
// try to re-establish connection
if (!connection.socket.writable) {
console.log("connection not writable - trying to reconnect");
await connection.connect(params).catch(error => {
console.log("connection error: ", error);
});
console.log("reconnect result: ", connection.socket.writable);
}
if ( connection.socket.writable === false ) {
res.status(500).json({message: 'telnet connection is not established.'});
try {
let response = await telnetPool.send('muc:room(\"'+ target +'\"):set_affiliation(true, \"'+ user +'\", \"' + affiliation +'\");\n');
if ((response.indexOf("nil value") > -1) || (response.indexOf("Fatal ") > -1)) {
res.status(500).json(response);
} else {
let response = await connection.send('muc:room(\"'+ target +'\"):set_affiliation(true, \"'+ user +'\", \"' + affiliation +'\");\n');
if ((response.indexOf("nil value") > -1) || (response.indexOf("Fatal ") > -1)) {
res.status(500).json(response);
} else {
let rawresults = response.split("\n");
let mucresults = [];
for (var i = 0; i < rawresults.length - 1 ; i++) {
if ((rawresults[i].startsWith("\u0000| ") || rawresults[i].startsWith("\| ")) && !rawresults[i].startsWith("| OK") && !rawresults[i].startsWith("\u0000| Result") ) {
console.log("adding: ", rawresults[i].split("| ")[1].split("\r")[0]);
mucresults.push(rawresults[i].split("| ")[1].split("\r")[0]);
}
let rawresults = response.split("\n");
let mucresults = [];
for (var i = 0; i < rawresults.length - 1 ; i++) {
if ((rawresults[i].startsWith("\u0000| ") || rawresults[i].startsWith("\| ")) && !rawresults[i].startsWith("| OK") && !rawresults[i].startsWith("\u0000| Result") ) {
console.log("adding: ", rawresults[i].split("| ")[1].split("\r")[0]);
mucresults.push(rawresults[i].split("| ")[1].split("\r")[0]);
}
console.log('async rawresult:', rawresults);
console.log('async result:', mucresults);
res.status(200).json(mucresults);
}
console.log('async rawresult:', rawresults);
console.log('async result:', mucresults);
res.status(200).json(mucresults);
}
} else {
res.status(404).json({message: 'connection not available'});
} catch (error) {
console.log("telnet send error: ", error);
res.status(500).json({message: 'telnet connection is not established.'});
}
} else {
// no groupchat found for requested operation - send 410 gone
@@ -499,52 +469,41 @@ router.post('/creategroup', async function (req, res) {
}
}
if (mucJidValid && req.body && req.body.jid && req.body.jid != "" && req.body.subject && req.body.subject != "" && req.body.owner && req.body.owner != "") {
if ( connection && connection.socket) {
if (!connection.socket.writable) {
await connection.connect(params).catch(error => {
console.log("connection error: ", error);
});
}
let alreadyExists = doesMucExist(req.body.jid);
console.log("existenceCheck ", alreadyExists);
if (alreadyExists == 1) {
res.status(409).json({message: "muc already exists"});
} else {
if (connection.socket.writable) {
console.log("creating room...");
let cmd = 'muc:create("' + req.body.jid.toLowerCase() + '", {';
cmd += 'subject="' + req.body.subject + '", history_length = 5, persistent = true });';
let response = await connection.send(cmd);
console.log(response);
if (response.indexOf("Result: MUC room (" + req.body.jid.toLowerCase()) > -1 ) {
// all good
console.log("setting owner to " + req.body.owner);
cmd = 'muc:room("' + req.body.jid + '"):set_affiliation(true, "' + req.body.owner + '", "owner");';
if (req.body.members && req.body.members.length > 0) {
for (var i = 0; i < req.body.members.length; i++) {
console.log("adding member " + req.body.members[i]);
cmd += 'muc:room("' + req.body.jid + '"):set_affiliation(true, "' + req.body.members[i] + '", "member");';
}
}
cmd += 'muc:room("' + req.body.jid + '"):save(true);';
response = await connection.send(cmd);
console.log(response);
//muc:room("testingtelnetcreate1_meeting@conference.microlab.zimbra-vnc.de"):set_affiliation(true, "richard.watson@microlab.zimbra-vnc.de", "member");
res.json({status: "ok"});
} else {
res.status(500).json({error: response});
}
} else {
console.log("connection not writable");
res.status(500).json({message: "connection not writable"});
}
}
let alreadyExists = await doesMucExist(req.body.jid);
console.log("existenceCheck ", alreadyExists);
if (alreadyExists == 1) {
res.status(409).json({message: "muc already exists"});
} else {
console.log("connection not writable2");
res.status(500).json({message: "connection not writable"});
try {
console.log("creating room...");
let cmd = 'muc:create("' + req.body.jid.toLowerCase() + '", {';
cmd += 'subject="' + req.body.subject + '", history_length = 5, persistent = true });';
let response = await telnetPool.send(cmd);
console.log(response);
if (response.indexOf("Result: MUC room (" + req.body.jid.toLowerCase()) > -1 ) {
// all good
console.log("setting owner to " + req.body.owner);
cmd = 'muc:room("' + req.body.jid + '"):set_affiliation(true, "' + req.body.owner + '", "owner");';
if (req.body.members && req.body.members.length > 0) {
for (var i = 0; i < req.body.members.length; i++) {
console.log("adding member " + req.body.members[i]);
cmd += 'muc:room("' + req.body.jid + '"):set_affiliation(true, "' + req.body.members[i] + '", "member");';
}
}
cmd += 'muc:room("' + req.body.jid + '"):save(true);';
response = await telnetPool.send(cmd);
console.log(response);
res.json({status: "ok"});
} else {
res.status(500).json({error: response});
}
} catch (error) {
console.log("telnet send error: ", error);
res.status(500).json({message: "connection not writable"});
}
}
} else {
res.status(500).json({message: "insufficient params"});
@@ -554,45 +513,35 @@ router.post('/creategroup', async function (req, res) {
router.get('/healthcheck', async function (req, res) {
if ( connection && connection.socket) {
if (!connection.socket.writable) {
await connection.connect(params).catch(error => {
console.log("connection error: ", error);
});
}
if (connection.socket.writable) {
console.log("getting c2s connections...");
let response = await connection.send('c2s:show();\n');
console.log(response);
let rawresults;
let resArr = [];
if (response.indexOf("https://prosody.im/doc/console\r\n") > -1) {
rawresults = response.split("https://prosody.im/doc/console\r\n")[1];
let rawDataArr = rawresults.split(" ");
for (var i = 0; i < rawDataArr.length - 1 ; i++) {
if (rawDataArr[i].indexOf("@") > -1) {
let jid = rawDataArr[i].split("/")[0];
resArr.push(jid);
}
try {
let response = await telnetPool.send('c2s:show();\n');
console.log(response);
let rawresults;
let resArr = [];
if (response.indexOf("https://prosody.im/doc/console\r\n") > -1) {
rawresults = response.split("https://prosody.im/doc/console\r\n")[1];
let rawDataArr = rawresults.split(" ");
for (var i = 0; i < rawDataArr.length - 1 ; i++) {
if (rawDataArr[i].indexOf("@") > -1) {
let jid = rawDataArr[i].split("/")[0];
resArr.push(jid);
}
console.log("processed: ", resArr);
} else {
rawresults = response;
let rawDataArr = rawresults.split(" ");
for (var i = 0; i < rawDataArr.length - 1 ; i++) {
if (rawDataArr[i].indexOf("@") > -1) {
let jid = rawDataArr[i].split("/")[0];
resArr.push(jid);
}
}
console.log("processed: ", resArr);
}
res.json({res: resArr});
console.log("processed: ", resArr);
} else {
console.log("connection not writable");
res.json({status: "ok"});
rawresults = response;
let rawDataArr = rawresults.split(" ");
for (var i = 0; i < rawDataArr.length - 1 ; i++) {
if (rawDataArr[i].indexOf("@") > -1) {
let jid = rawDataArr[i].split("/")[0];
resArr.push(jid);
}
}
console.log("processed: ", resArr);
}
} else {
res.json({res: resArr});
} catch (error) {
console.log("telnet send error: ", error);
res.json({status: "ok"});
}
});
+99
View File
@@ -0,0 +1,99 @@
var Telnet = require('telnet-client');
var STATE = { IDLE: 0, BUSY: 1, DEAD: 2 };
function TelnetPool(params, size) {
this.params = params;
this.size = size;
this.conns = [];
this.queue = [];
for (var i = 0; i < size; i++) {
this.conns.push(this._spawn('#' + i));
}
}
TelnetPool.prototype._spawn = function (label) {
var self = this;
var conn = new Telnet();
conn._poolState = STATE.DEAD;
conn._poolLabel = label;
conn.connect(self.params).then(function () {
conn._poolState = STATE.IDLE;
console.log('[telnet-pool] connected', label);
self._drain();
}).catch(function (error) {
conn._poolState = STATE.DEAD;
console.log('[telnet-pool] connect error', label, error);
});
return conn;
};
TelnetPool.prototype._drain = function () {
if (this.queue.length === 0) return;
for (var i = 0; i < this.conns.length; i++) {
var conn = this.conns[i];
if (conn._poolState === STATE.IDLE) {
conn._poolState = STATE.BUSY;
var resolve = this.queue.shift();
resolve(conn);
if (this.queue.length === 0) return;
}
}
};
TelnetPool.prototype._acquire = function () {
var self = this;
for (var i = 0; i < self.conns.length; i++) {
var conn = self.conns[i];
if (conn._poolState === STATE.IDLE) {
conn._poolState = STATE.BUSY;
return Promise.resolve(conn);
}
}
return new Promise(function (resolve) {
self.queue.push(resolve);
});
};
TelnetPool.prototype._release = function (conn) {
conn._poolState = STATE.IDLE;
this._drain();
};
TelnetPool.prototype._ensureConnected = function (conn) {
var self = this;
if (conn && conn.socket && conn.socket.writable) {
return Promise.resolve(conn);
}
console.log('[telnet-pool] socket not writable, reconnecting', conn._poolLabel);
return conn.connect(self.params).then(function () {
console.log('[telnet-pool] reconnected', conn._poolLabel);
return conn;
});
};
TelnetPool.prototype._replace = function (deadConn) {
var idx = this.conns.indexOf(deadConn);
if (idx > -1) {
this.conns[idx] = this._spawn(deadConn._poolLabel);
}
};
TelnetPool.prototype.send = function (cmd) {
var self = this;
return self._acquire().then(function (conn) {
return self._ensureConnected(conn).then(function () {
return conn.send(cmd);
}).then(function (response) {
self._release(conn);
return response;
}).catch(function (error) {
self._replace(conn);
throw error;
});
});
};
module.exports = function (params, size) {
return new TelnetPool(params, size);
};
+2 -1
View File
@@ -9,7 +9,8 @@ module.exports = {
host: process.env.telnetHost || "127.0.0.1",
port: process.env.telneetPort || 5582,
negotiationMandatory: false,
timeout: process.env.telnetTimeout || 1500
timeout: process.env.telnetTimeout || 5000,
poolSize: parseInt(process.env.telnetPoolSize, 10) || 4
},
database: {
host: process.env.dbHost || "127.0.0.1",