From 7f129a78bf07049757db15e52bfa5e1bc9fc4652 Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Fri, 27 Mar 2026 13:10:19 +0200 Subject: [PATCH 01/16] fix online list request --- APIs/JSON/routes/packet_processor.js | 8 ++++---- app/providers/services/session/index.js | 12 ++++++++---- 2 files changed, 12 insertions(+), 8 deletions(-) diff --git a/APIs/JSON/routes/packet_processor.js b/APIs/JSON/routes/packet_processor.js index 2a473bd7..ca93bd78 100644 --- a/APIs/JSON/routes/packet_processor.js +++ b/APIs/JSON/routes/packet_processor.js @@ -46,18 +46,18 @@ class PacketJsonProcessor extends BasePacketProcessor { } catch (error) { logger.error(error) let errorBackMessage = null - if (json.request) { + if (json?.request) { errorBackMessage = { response: { - id: json.request.id, + id: json.request?.id, error: error.cause || error.message, }, } } else { - const topLevelElement = Object.keys(json)[0] + const topLevelElement = json ? Object.keys(json)[0] : void 0 errorBackMessage = { [topLevelElement]: { - id: json[topLevelElement].id, + id: json?.[topLevelElement]?.id, error: error.cause || error.message, }, } diff --git a/app/providers/services/session/index.js b/app/providers/services/session/index.js index 55f0ce29..1bd0b406 100644 --- a/app/providers/services/session/index.js +++ b/app/providers/services/session/index.js @@ -102,9 +102,7 @@ class SessionService { async listUserDevice(organizationId, userId) { if (this.config.get("app.isStandAloneNode")) { - return this.getUserDevices(userId) - .map((connection) => connection?.deviceId) - .filter((deviceId) => deviceId !== CONSTANTS.HTTP_DEVICE_ID) + return this.listUserDeviceLocal(userId) } const userKey = this.#usersSetCacheKey(organizationId, userId) @@ -113,6 +111,12 @@ class SessionService { return deviceIds ?? [] } + listUserDeviceLocal(userId) { + return this.getUserDevices(userId) + .map((connection) => connection?.deviceId) + .filter((deviceId) => deviceId !== CONSTANTS.HTTP_DEVICE_ID) + } + async deleteUserDevices(organizationId, userId) { const userKey = this.#usersSetCacheKey(organizationId, userId) @@ -393,7 +397,7 @@ class SessionService { (session) => session?.organizationId === organizationId && session?.extraParams[CONSTANTS.SESSION_DEVICE_ID_KEY] !== CONSTANTS.HTTP_DEVICE_ID && - session?.userId + session?.userId && this.listUserDeviceLocal(session?.userId)?.length ) .map((session) => session.userId) .sort((userIdA, userIdB) => userIdA - userIdB) From ab0007d0eafe8953bf899e8e825bd780cc921700 Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Fri, 17 Apr 2026 14:16:42 +0300 Subject: [PATCH 02/16] add basic cmd ws test commands --- index.js | 56 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 56 insertions(+) diff --git a/index.js b/index.js index fee58008..41a35bab 100644 --- a/index.js +++ b/index.js @@ -195,4 +195,60 @@ if (config.get("tcp.isEnabled")) { await tcpProtocolImp.listen(tcpOptions) } +process.stdin.setEncoding('utf8') +process.stdin.on('data', (data) => { + try { + const cmd = data.trim() + console.log('[Cmd]', cmd) + + const sessionService = ServiceLocatorContainer.use("SessionService") + // const findSocketByUserId = userId => { + // for (const socket of sessionService.activeSessions.SESSIONS.keys()) { + // const userData = sessionService.activeSessions.SESSIONS.get(socket) + + // if (userData?.userId === userId) { + // return socket + // } + // } + // } + + if (cmd.match(/cmd-ping/i)) { // 'cmd-ping 1111' + const matchRes = cmd.match(/cmd-ping (.+)/i) + const userId = +matchRes.at(1) + + console.log('[PingWS]', userId, '[devices]', sessionService.listUserDeviceLocal(userId)) + + const connections = sessionService.listUserDeviceLocal(userId) + + for (const connection of connections) { + console.log('[PingWS][start]', connection?.socket) + + const sendResult = connection?.socket?.ping(" ") + + console.log('[PingWS][result]', connection?.socket, sendResult) + } + } + + if (cmd.match(/cmd-send/i)) { // 'cmd-send 1111 test' + const matchRes = cmd.match(/cmd-send (.+) (.+)/i) + const userId = +matchRes.at(1) + const sendData = matchRes.at(2) + + console.log('[SendWS]', userId, sendData, '[devices]', sessionService.listUserDeviceLocal(userId)) + + const connections = sessionService.listUserDeviceLocal(userId) + + for (const connection of connections) { + console.log('[SendWS][start]', connection?.socket) + + const sendResult = connection?.socket?.send(sendData) + + console.log('[SendWS][result]', connection?.socket, sendResult) + } + } + } catch (error) { + console.log('[Cmd][error]', error) + } +}) + // https://dev.to/mattkrick/replacing-express-with-uwebsockets-48ph From e16949d15e3ea63c235e555cdc3c0a969f31534c Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Fri, 17 Apr 2026 15:19:33 +0300 Subject: [PATCH 03/16] update check cmd connections --- index.js | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/index.js b/index.js index 41a35bab..b38fdf26 100644 --- a/index.js +++ b/index.js @@ -218,7 +218,7 @@ process.stdin.on('data', (data) => { console.log('[PingWS]', userId, '[devices]', sessionService.listUserDeviceLocal(userId)) - const connections = sessionService.listUserDeviceLocal(userId) + const connections = sessionService.getUserDevices(userId) for (const connection of connections) { console.log('[PingWS][start]', connection?.socket) @@ -236,7 +236,7 @@ process.stdin.on('data', (data) => { console.log('[SendWS]', userId, sendData, '[devices]', sessionService.listUserDeviceLocal(userId)) - const connections = sessionService.listUserDeviceLocal(userId) + const connections = sessionService.getUserDevices(userId) for (const connection of connections) { console.log('[SendWS][start]', connection?.socket) From d6e4c9b104fb1f83c7dc83f01c757117cc4930ac Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Fri, 17 Apr 2026 16:44:51 +0300 Subject: [PATCH 04/16] add close --- index.js | 19 ++++++++++++++++++- 1 file changed, 18 insertions(+), 1 deletion(-) diff --git a/index.js b/index.js index b38fdf26..91452ab2 100644 --- a/index.js +++ b/index.js @@ -230,7 +230,7 @@ process.stdin.on('data', (data) => { } if (cmd.match(/cmd-send/i)) { // 'cmd-send 1111 test' - const matchRes = cmd.match(/cmd-send (.+) (.+)/i) + const matchRes = cmd.match(/cmd-send ([^\s]+) (.+)/i) const userId = +matchRes.at(1) const sendData = matchRes.at(2) @@ -246,6 +246,23 @@ process.stdin.on('data', (data) => { console.log('[SendWS][result]', connection?.socket, sendResult) } } + + if (cmd.match(/cmd-close/i)) { // 'cmd-close 1111' + const matchRes = cmd.match(/cmd-close (.+)/i) + const userId = +matchRes.at(1) + + console.log('[CloseWS]', userId, '[devices]', sessionService.listUserDeviceLocal(userId)) + + const connections = sessionService.getUserDevices(userId) + + for (const connection of connections) { + console.log('[CloseWS][start]', connection?.socket) + + const sendResult = connection?.socket?.close() + + console.log('[CloseWS][result]', connection?.socket, sendResult) + } + } } catch (error) { console.log('[Cmd][error]', error) } From 83abc3942d7d731bece7b44dd0d99503a454ca09 Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Tue, 5 May 2026 12:39:11 +0300 Subject: [PATCH 05/16] add socketCloseWatchdog --- app/config/default.js | 1 + app/utils/socket-close-watchdog.js | 28 ++++++++++++++++++++ index.js | 41 +++++++++++++++++++++++++----- 3 files changed, 64 insertions(+), 6 deletions(-) create mode 100644 app/utils/socket-close-watchdog.js diff --git a/app/config/default.js b/app/config/default.js index 5ea61b78..3e85bb89 100644 --- a/app/config/default.js +++ b/app/config/default.js @@ -7,6 +7,7 @@ const CONFIG = { name: process.env.APP_NAME ?? "SAMA", hostName: process.env.HOSTNAME, isStandAloneNode: process.env.STANDALONE_NODE === CONSTANTS.ENV_TRUE, + socketCloseWatchdogInterval: +process.env.SOCKET_CLOSE_WATCHDOG_INTERVAL }, logger: { logLevel: process.env.LOG_LEVEL ?? "debug", diff --git a/app/utils/socket-close-watchdog.js b/app/utils/socket-close-watchdog.js new file mode 100644 index 00000000..b8dc5f40 --- /dev/null +++ b/app/utils/socket-close-watchdog.js @@ -0,0 +1,28 @@ +import net from "node:net" + +export const socketCloseWatchdog = (logger, sessionService, onWsCloseCb, onTcpCloseCb) => { + const users = Object.keys(sessionService.activeSessions.DEVICES) + + logger.debug("[run] %s", users.length) + + for (const userId of users) { + const connections = sessionService.activeSessions.DEVICES[userId] ?? [] + for (const connection of connections) { + const isTCP = connection?.socket instanceof net.Socket + try { + if (isTCP) { + connection?.socket?.write(" ") + } else { + connection?.socket?.send(" ") + } + } catch (error) { + logger.error(error, "[error socket send] %s", userId) + if (isTCP) { + onTcpCloseCb(connection?.socket) + } else { + onWsCloseCb(connection?.socket, 10) + } + } + } + } +} \ No newline at end of file diff --git a/index.js b/index.js index 91452ab2..7e36ec28 100644 --- a/index.js +++ b/index.js @@ -24,6 +24,7 @@ import OTPSender from "./app/lib/otp_sender.js" import { APIs } from "./app/networking/APIs.js" import { buildWsEndpoint } from "./app/utils/build_ws_endpoint.js" +import { socketCloseWatchdog } from "./app/utils/socket-close-watchdog.js" if (config.get("app.env") === CONSTANTS.ENVS.PROD) { process.on("unhandledRejection", (reason, promise) => { @@ -190,11 +191,26 @@ await wsProtocolImp.listen(uWSOptions) const httpProtocolImp = new HttpProtocol(sessionService, conversationService, wsProtocolImp.uWSocketServer) await httpProtocolImp.listen({}) +// https://dev.to/mattkrick/replacing-express-with-uwebsockets-48ph + +let tcpProtocolImp = void 0 if (config.get("tcp.isEnabled")) { - const tcpProtocolImp = new TcpProtocol(sessionService, conversationService) + tcpProtocolImp = new TcpProtocol(sessionService, conversationService) await tcpProtocolImp.listen(tcpOptions) } +if (config.get("app.socketCloseWatchdogInterval")) { + const socketCloseWatchdogLogger = logger.child("[SocketClosedWatchDog]") + setInterval(() => { + socketCloseWatchdog( + socketCloseWatchdogLogger, + sessionService, + (socket, code) => wsProtocolImp.onClose(socket, code), + (socket) => tcpProtocolImp.onClose(socket) + ) + }, config.get("app.socketCloseWatchdogInterval")) +} + process.stdin.setEncoding('utf8') process.stdin.on('data', (data) => { try { @@ -223,7 +239,12 @@ process.stdin.on('data', (data) => { for (const connection of connections) { console.log('[PingWS][start]', connection?.socket) - const sendResult = connection?.socket?.ping(" ") + let sendResult = void 0 + try { + sendResult = connection?.socket?.ping(" ") + } catch (err) { + console.log('[Cmd][error]', err) + } console.log('[PingWS][result]', connection?.socket, sendResult) } @@ -241,7 +262,12 @@ process.stdin.on('data', (data) => { for (const connection of connections) { console.log('[SendWS][start]', connection?.socket) - const sendResult = connection?.socket?.send(sendData) + let sendResult = void 0 + try { + sendResult = connection?.socket?.send(sendData) + } catch (err) { + console.log('[Cmd][error]', err) + } console.log('[SendWS][result]', connection?.socket, sendResult) } @@ -258,7 +284,12 @@ process.stdin.on('data', (data) => { for (const connection of connections) { console.log('[CloseWS][start]', connection?.socket) - const sendResult = connection?.socket?.close() + let sendResult = void 0 + try { + sendResult = connection?.socket?.close() + } catch (err) { + console.log('[Cmd][error]', err) + } console.log('[CloseWS][result]', connection?.socket, sendResult) } @@ -267,5 +298,3 @@ process.stdin.on('data', (data) => { console.log('[Cmd][error]', error) } }) - -// https://dev.to/mattkrick/replacing-express-with-uwebsockets-48ph From bbc7b851c009dcf66c1f1755d903f5a9cf58fce2 Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Tue, 5 May 2026 12:56:49 +0300 Subject: [PATCH 06/16] add remove session --- app/utils/socket-close-watchdog.js | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/app/utils/socket-close-watchdog.js b/app/utils/socket-close-watchdog.js index b8dc5f40..403c8f28 100644 --- a/app/utils/socket-close-watchdog.js +++ b/app/utils/socket-close-watchdog.js @@ -1,6 +1,6 @@ import net from "node:net" -export const socketCloseWatchdog = (logger, sessionService, onWsCloseCb, onTcpCloseCb) => { +export const socketCloseWatchdog = async (logger, sessionService, onWsCloseCb, onTcpCloseCb) => { const users = Object.keys(sessionService.activeSessions.DEVICES) logger.debug("[run] %s", users.length) @@ -18,10 +18,11 @@ export const socketCloseWatchdog = (logger, sessionService, onWsCloseCb, onTcpCl } catch (error) { logger.error(error, "[error socket send] %s", userId) if (isTCP) { - onTcpCloseCb(connection?.socket) + await onTcpCloseCb(connection?.socket).catch(error => logger.error(error, "[close tcp]")) } else { - onWsCloseCb(connection?.socket, 10) + await onWsCloseCb(connection?.socket, 10).catch(error => logger.error(error, "[close ws]")) } + await sessionService.removeUserSession(connection?.socket, userId, connection?.deviceId).catch(error => logger.error(error, "[remove]")) } } } From e9fd59abaa18819243fa4db7bec51f2cd7a29ebc Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Tue, 5 May 2026 13:10:25 +0300 Subject: [PATCH 07/16] on close logs debug --- app/networking/protocol_processors/base.js | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/app/networking/protocol_processors/base.js b/app/networking/protocol_processors/base.js index ca9abf0c..66663ab7 100644 --- a/app/networking/protocol_processors/base.js +++ b/app/networking/protocol_processors/base.js @@ -200,7 +200,7 @@ class BaseProtocolProcessor { } async onClose(socket, code) { - logger.trace("[Close] IP: %s CLIENT_ID: %s CODE: %s", this.socketAddress(socket), socket.clientId, code) + logger.debug("[Close] IP: %s CLIENT_ID: %s CODE: %s", this.socketAddress(socket), socket.clientId, code) socket.isAlive = false @@ -214,7 +214,7 @@ class BaseProtocolProcessor { async updateLastUserLastActivityOnClose(socket) { const { organizationId, userId } = this.sessionService.getSession(socket) ?? {} - logger.trace("[UPDATE_LAST_ACTIVITY][CLOSE] OrgId: %s UserId: %s", organizationId, userId) + logger.debug("[UPDATE_LAST_ACTIVITY][CLOSE] OrgId: %s UserId: %s", organizationId, userId) if (!userId) { return From babfcd82bb5bf0f7e5872435bfd01f5fef48bce0 Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Tue, 5 May 2026 13:20:39 +0300 Subject: [PATCH 08/16] add repl service --- .env.example | 10 +++ app/config/default.js | 13 ++++ app/lib/repl-tools.js | 144 ++++++++++++++++++++++++++++++++++++++++++ index.js | 94 ++------------------------- 4 files changed, 174 insertions(+), 87 deletions(-) create mode 100644 app/lib/repl-tools.js diff --git a/.env.example b/.env.example index e51f056b..f6bb36ea 100644 --- a/.env.example +++ b/.env.example @@ -11,6 +11,16 @@ MONGODB_URL=mongodb://127.0.0.1/samadb REDIS_URL=redis://127.0.0.1:6379 +# REPL +# http +APP_REPL_HTTP_PORT=5010 +APP_REPL_HTTP_ACCESS_KEY=repl-wow-key +# socket +APP_REPL_SOCKET_HANDLER=/tmp/net-repl.socket +# file +APP_REPL_FILE_IN=/tmp/pipe-repl.in +APP_REPL_FILE_OUT=/tmp/pipe-repl.out + STORAGE_DRIVER=minio # If you set STORAGE_DRIVER=s3, then fill the below envs diff --git a/app/config/default.js b/app/config/default.js index 3e85bb89..6621c5ac 100644 --- a/app/config/default.js +++ b/app/config/default.js @@ -103,6 +103,19 @@ const CONFIG = { apiKey: process.env.HTTP_ADMIN_API_KEY, }, }, + repl: { + http: { + port: +(process.env.APP_REPL_HTTP_PORT ?? 5010), + accessKey: process.env.APP_REPL_HTTP_ACCESS_KEY, + }, + socket: { + handler: process.env.APP_REPL_SOCKET_HANDLER, + }, + file: { + in: process.env.APP_REPL_FILE_IN, + out: process.env.APP_REPL_FILE_OUT, + }, + }, conversation: { disableChannelsLogic: process.env.CONVERSATION_DISABLE_CHANNELS_LOGIC === CONSTANTS.ENV_TRUE, isEventsEnabled: process.env.CONVERSATION_NOTIFICATIONS_ENABLED === CONSTANTS.ENV_TRUE, diff --git a/app/lib/repl-tools.js b/app/lib/repl-tools.js new file mode 100644 index 00000000..186a9b96 --- /dev/null +++ b/app/lib/repl-tools.js @@ -0,0 +1,144 @@ +import fs from "node:fs" +import fsPromises from "node:fs/promises" +import vm from "node:vm" +import { exec } from "node:child_process" +import repl from "node:repl" +import net from "node:net" +import http from "node:http" + +import { CONSTANTS } from "../constants/constants.js" + +// example: curl --no-progress-meter -sSNT. -H "api-key: ****" localhost:5010 + +const httpReplService = async (replOptions, httpOptions) => { + const { ctx } = replOptions + const { accessKey, port } = httpOptions + + const server = http.createServer((req, res) => { + if (req.headers[CONSTANTS.HTTP_REPL_ACCESS_KEY_HEADER] !== accessKey) { + return res.end("Invalid access-key") + } + + res.setHeader("content-type", "multipart/octet-stream") + + const replService = repl.start({ + prompt: "curl repl> ", + input: req, + output: res, + terminal: false, + useColors: true, + useGlobal: false, + }) + + const context = vm.createContext({ ...ctx }) + + replService.context = context + + req.on("error", (error) => replService.close()) + req.on("end", () => replService.close()) + + res.on("error", (error) => replService.close()) + res.on("end", () => replService.close()) + + replService.on("close", () => !res.closed && res.end("REPL closed")) + replService.on("error", (error) => { + replService.close() + !res.closed && res.end("REPL closed") + !req.closed && req.destroy() + }) + }) + + return new Promise((resolve, reject) => { + server.listen(port, () => resolve(port)) + }) +} + +// example write: cat > ./pipe-repl.in +// example read: tail -f ./pipe-repl.out + +const fileReplService = async (replOptions, fileOptions) => { + const { ctx } = replOptions + const { fileIn, fileOut } = fileOptions + + await fsPromises.rm(fileIn, { force: true }).catch((error) => {}) + await fsPromises.rm(fileOut, { force: true }).catch((error) => {}) + + await new Promise((resolve, reject) => { + exec(`mkfifo ${fileIn}`, (error, out, outErr) => { + if (error) { + return reject(error) + } + resolve() + }) + }) + + const cmdInputPipe = fs.createReadStream(fileIn, { encoding: "utf8" }) + const cmdOutPipe = fs.createWriteStream(fileOut, { encoding: "utf8" }) + + const replService = repl.start({ + prompt: "file repl> ", + input: cmdInputPipe, + output: cmdOutPipe, + terminal: false, + useColors: true, + useGlobal: false, + }) + + const context = vm.createContext({ ...ctx }) + + replService.context = context + + replService.on("error", (error) => { + replService.close() + cmdInputPipe.close() + }) + + cmdInputPipe.on("end", () => { + replService.close() + fileReplService(replOptions, fileOptions) + }) +} + +// example: nc -U ./net-repl.socket + +const netReplService = async (replOptions, fileOptions) => { + const { ctx } = replOptions + const { socketHandler } = fileOptions + + await fsPromises.rm(socketHandler, { force: true }).catch((error) => {}) + + const server = net.createServer((socket) => { + const replService = repl.start({ + prompt: "socket repl> ", + input: socket, + output: socket, + terminal: false, + useColors: true, + useGlobal: false, + }) + + const context = vm.createContext({ ...ctx }) + + replService.context = context + + socket.on("end", () => replService.close()) + socket.on("error", (error) => replService.close()) + replService.on("close", () => !socket.closed && socket.end("REPL closed")) + }) + + return new Promise((resolve, reject) => { + server.listen(socketHandler, () => resolve(socketHandler)) + }) +} + +export const startReplServices = async (replOptions, httpOptions, netOptions, fileOptions) => { + if (httpOptions.accessKey) { + await httpReplService(replOptions, httpOptions) + } + if (netOptions.socketHandler) { + await netReplService(replOptions, netOptions) + } + if (fileOptions.fileIn && fileOptions.fileOut) { + await fileReplService(replOptions, fileOptions) + } +} \ No newline at end of file diff --git a/index.js b/index.js index 7e36ec28..c3782290 100644 --- a/index.js +++ b/index.js @@ -20,6 +20,7 @@ import HttpProtocol from "./app/networking/protocol_processors/http.js" import { connectToDBPromise } from "./app/lib/db.js" import RedisClient from "./app/lib/redis.js" import OTPSender from "./app/lib/otp_sender.js" +import { startReplServices } from "./app/lib/repl-tools.js" import { APIs } from "./app/networking/APIs.js" @@ -211,90 +212,9 @@ if (config.get("app.socketCloseWatchdogInterval")) { }, config.get("app.socketCloseWatchdogInterval")) } -process.stdin.setEncoding('utf8') -process.stdin.on('data', (data) => { - try { - const cmd = data.trim() - console.log('[Cmd]', cmd) - - const sessionService = ServiceLocatorContainer.use("SessionService") - // const findSocketByUserId = userId => { - // for (const socket of sessionService.activeSessions.SESSIONS.keys()) { - // const userData = sessionService.activeSessions.SESSIONS.get(socket) - - // if (userData?.userId === userId) { - // return socket - // } - // } - // } - - if (cmd.match(/cmd-ping/i)) { // 'cmd-ping 1111' - const matchRes = cmd.match(/cmd-ping (.+)/i) - const userId = +matchRes.at(1) - - console.log('[PingWS]', userId, '[devices]', sessionService.listUserDeviceLocal(userId)) - - const connections = sessionService.getUserDevices(userId) - - for (const connection of connections) { - console.log('[PingWS][start]', connection?.socket) - - let sendResult = void 0 - try { - sendResult = connection?.socket?.ping(" ") - } catch (err) { - console.log('[Cmd][error]', err) - } - - console.log('[PingWS][result]', connection?.socket, sendResult) - } - } - - if (cmd.match(/cmd-send/i)) { // 'cmd-send 1111 test' - const matchRes = cmd.match(/cmd-send ([^\s]+) (.+)/i) - const userId = +matchRes.at(1) - const sendData = matchRes.at(2) - - console.log('[SendWS]', userId, sendData, '[devices]', sessionService.listUserDeviceLocal(userId)) - - const connections = sessionService.getUserDevices(userId) - - for (const connection of connections) { - console.log('[SendWS][start]', connection?.socket) - - let sendResult = void 0 - try { - sendResult = connection?.socket?.send(sendData) - } catch (err) { - console.log('[Cmd][error]', err) - } - - console.log('[SendWS][result]', connection?.socket, sendResult) - } - } - - if (cmd.match(/cmd-close/i)) { // 'cmd-close 1111' - const matchRes = cmd.match(/cmd-close (.+)/i) - const userId = +matchRes.at(1) - - console.log('[CloseWS]', userId, '[devices]', sessionService.listUserDeviceLocal(userId)) - - const connections = sessionService.getUserDevices(userId) - - for (const connection of connections) { - console.log('[CloseWS][start]', connection?.socket) - - let sendResult = void 0 - try { - sendResult = connection?.socket?.close() - } catch (err) { - console.log('[Cmd][error]', err) - } - - console.log('[CloseWS][result]', connection?.socket, sendResult) - } - } - } catch (error) { - console.log('[Cmd][error]', error) - } -}) +await startReplServices( + { ctx: { slc: ServiceLocatorContainer } }, + { accessKey: config.get("repl.http.accessKey"), port: config.get("repl.http.port") }, + { socketHandler: config.get("repl.socket.handler") }, + { fileIn: config.get("repl.file.in"), fileOut: config.get("repl.file.out") } +) \ No newline at end of file From 3b63d0d4c6d4f4d79f6ca592004a8efbbf99198f Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Tue, 5 May 2026 13:30:19 +0300 Subject: [PATCH 09/16] update docker-file --- Dockerfile | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Dockerfile b/Dockerfile index 0dc79b73..ec87285b 100644 --- a/Dockerfile +++ b/Dockerfile @@ -11,7 +11,7 @@ COPY . . FROM node:22-slim RUN apt-get update && apt-get install -y --no-install-recommends \ - curl \ + curl netcat-openbsd \ && rm -rf /var/lib/apt/lists/* WORKDIR /app From 33844eff518c53a434376da1ea87097bb9cf8bd9 Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Tue, 5 May 2026 14:33:02 +0300 Subject: [PATCH 10/16] add ping socket --- .env.example | 1 + APIs/JSON/index.js | 6 ++++++ APIs/XMPP | 2 +- app/config/default.js | 2 +- ...-close-watchdog.js => watchdog-ping-socket.js} | 15 +++++++++++---- index.js | 10 +++++----- 6 files changed, 25 insertions(+), 11 deletions(-) rename app/utils/{socket-close-watchdog.js => watchdog-ping-socket.js} (65%) diff --git a/.env.example b/.env.example index f6bb36ea..7c50e63a 100644 --- a/.env.example +++ b/.env.example @@ -1,5 +1,6 @@ NODE_ENV=development STANDALONE_NODE=false +WATCHDOG_PING_SOCKET_INTERVAL=60000 # 1 minute APP_PORT=9001 APP_TCP_PORT=8001 diff --git a/APIs/JSON/index.js b/APIs/JSON/index.js index d40c2ece..55f4071c 100644 --- a/APIs/JSON/index.js +++ b/APIs/JSON/index.js @@ -35,4 +35,10 @@ export default class JsonAPI extends BaseAPI { } return this.stringifyMessage(message) } + + pingPackage() { + const ping = { response: { ping: {} } } + + return this.stringifyMessage(ping) + } } diff --git a/APIs/XMPP b/APIs/XMPP index ba4c3af7..c16b8542 160000 --- a/APIs/XMPP +++ b/APIs/XMPP @@ -1 +1 @@ -Subproject commit ba4c3af760d78fb9da77e64949e6186b99878f0f +Subproject commit c16b85422ea0520b1e422a954642b6bce4ccac96 diff --git a/app/config/default.js b/app/config/default.js index 6621c5ac..c73ac88e 100644 --- a/app/config/default.js +++ b/app/config/default.js @@ -7,7 +7,7 @@ const CONFIG = { name: process.env.APP_NAME ?? "SAMA", hostName: process.env.HOSTNAME, isStandAloneNode: process.env.STANDALONE_NODE === CONSTANTS.ENV_TRUE, - socketCloseWatchdogInterval: +process.env.SOCKET_CLOSE_WATCHDOG_INTERVAL + watchdogPingSocketInterval: +process.env.WATCHDOG_PING_SOCKET_INTERVAL, }, logger: { logLevel: process.env.LOG_LEVEL ?? "debug", diff --git a/app/utils/socket-close-watchdog.js b/app/utils/watchdog-ping-socket.js similarity index 65% rename from app/utils/socket-close-watchdog.js rename to app/utils/watchdog-ping-socket.js index 403c8f28..d5c51aa7 100644 --- a/app/utils/socket-close-watchdog.js +++ b/app/utils/watchdog-ping-socket.js @@ -1,6 +1,7 @@ import net from "node:net" +import { APIs, BASE_API } from "../networking/APIs.js" -export const socketCloseWatchdog = async (logger, sessionService, onWsCloseCb, onTcpCloseCb) => { +export const watchdogPingSocket = async (logger, sessionService, onWsCloseCb, onTcpCloseCb) => { const users = Object.keys(sessionService.activeSessions.DEVICES) logger.debug("[run] %s", users.length) @@ -8,12 +9,18 @@ export const socketCloseWatchdog = async (logger, sessionService, onWsCloseCb, o for (const userId of users) { const connections = sessionService.activeSessions.DEVICES[userId] ?? [] for (const connection of connections) { - const isTCP = connection?.socket instanceof net.Socket + if (!connection?.socket) { + continue + } + + const isTCP = connection.socket instanceof net.Socket + const pingPackage = APIs[connection.socket.apiType ?? BASE_API].pingPackage() + try { if (isTCP) { - connection?.socket?.write(" ") + connection?.socket?.write(pingPackage) } else { - connection?.socket?.send(" ") + connection?.socket?.send(pingPackage) } } catch (error) { logger.error(error, "[error socket send] %s", userId) diff --git a/index.js b/index.js index c3782290..ca0b0de2 100644 --- a/index.js +++ b/index.js @@ -25,7 +25,7 @@ import { startReplServices } from "./app/lib/repl-tools.js" import { APIs } from "./app/networking/APIs.js" import { buildWsEndpoint } from "./app/utils/build_ws_endpoint.js" -import { socketCloseWatchdog } from "./app/utils/socket-close-watchdog.js" +import { watchdogPingSocket } from "./app/utils/watchdog-ping-socket.js" if (config.get("app.env") === CONSTANTS.ENVS.PROD) { process.on("unhandledRejection", (reason, promise) => { @@ -200,16 +200,16 @@ if (config.get("tcp.isEnabled")) { await tcpProtocolImp.listen(tcpOptions) } -if (config.get("app.socketCloseWatchdogInterval")) { +if (config.get("app.watchdogPingSocketInterval")) { const socketCloseWatchdogLogger = logger.child("[SocketClosedWatchDog]") setInterval(() => { - socketCloseWatchdog( + watchdogPingSocket( socketCloseWatchdogLogger, sessionService, (socket, code) => wsProtocolImp.onClose(socket, code), (socket) => tcpProtocolImp.onClose(socket) ) - }, config.get("app.socketCloseWatchdogInterval")) + }, config.get("app.watchdogPingSocketInterval")) } await startReplServices( @@ -217,4 +217,4 @@ await startReplServices( { accessKey: config.get("repl.http.accessKey"), port: config.get("repl.http.port") }, { socketHandler: config.get("repl.socket.handler") }, { fileIn: config.get("repl.file.in"), fileOut: config.get("repl.file.out") } -) \ No newline at end of file +) From 131dba8942752c4596a4b7f1d11259c7b4bf48b0 Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Tue, 5 May 2026 14:42:01 +0300 Subject: [PATCH 11/16] update package --- APIs/XMPP | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/APIs/XMPP b/APIs/XMPP index c16b8542..80ab3a71 160000 --- a/APIs/XMPP +++ b/APIs/XMPP @@ -1 +1 @@ -Subproject commit c16b85422ea0520b1e422a954642b6bce4ccac96 +Subproject commit 80ab3a714e5b67ac9a0c130d2e0c651bff82c381 From 4ef8713689ab679ca9596e0dcc64c528e5c95f79 Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Tue, 5 May 2026 15:17:48 +0300 Subject: [PATCH 12/16] fix tcp send failed --- app/utils/watchdog-ping-socket.js | 33 +++++++++++++++++++++++-------- 1 file changed, 25 insertions(+), 8 deletions(-) diff --git a/app/utils/watchdog-ping-socket.js b/app/utils/watchdog-ping-socket.js index d5c51aa7..5f98a9ba 100644 --- a/app/utils/watchdog-ping-socket.js +++ b/app/utils/watchdog-ping-socket.js @@ -1,6 +1,19 @@ import net from "node:net" import { APIs, BASE_API } from "../networking/APIs.js" +const onSentFailed = async ( + error, isTCP, userId, connection, + logger, sessionService, onTcpCloseCb, onWsCloseCb +) => { + logger.error(error, "[error socket send] %s", userId) + if (isTCP) { + await onTcpCloseCb(connection?.socket).catch(error => logger.error(error, "[close tcp]")) + } else { + await onWsCloseCb(connection?.socket, 10).catch(error => logger.error(error, "[close ws]")) + } + await sessionService.removeUserSession(connection?.socket, userId, connection?.deviceId).catch(error => logger.error(error, "[remove]")) +} + export const watchdogPingSocket = async (logger, sessionService, onWsCloseCb, onTcpCloseCb) => { const users = Object.keys(sessionService.activeSessions.DEVICES) @@ -18,18 +31,22 @@ export const watchdogPingSocket = async (logger, sessionService, onWsCloseCb, on try { if (isTCP) { - connection?.socket?.write(pingPackage) + connection?.socket?.write(pingPackage, (error) => { + if (error) { + onSentFailed( + error, isTCP, userId, connection, + logger, sessionService, onWsCloseCb, onTcpCloseCb + ) + } + }) } else { connection?.socket?.send(pingPackage) } } catch (error) { - logger.error(error, "[error socket send] %s", userId) - if (isTCP) { - await onTcpCloseCb(connection?.socket).catch(error => logger.error(error, "[close tcp]")) - } else { - await onWsCloseCb(connection?.socket, 10).catch(error => logger.error(error, "[close ws]")) - } - await sessionService.removeUserSession(connection?.socket, userId, connection?.deviceId).catch(error => logger.error(error, "[remove]")) + await onSentFailed( + error, isTCP, userId, connection, + logger, sessionService, onWsCloseCb, onTcpCloseCb + ) } } } From 616f35fce74c762dd6e37383f8035506368a65ea Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Tue, 5 May 2026 15:27:28 +0300 Subject: [PATCH 13/16] update watchdogPingSocket --- app/utils/watchdog-ping-socket.js | 41 +++++++++++-------------------- 1 file changed, 15 insertions(+), 26 deletions(-) diff --git a/app/utils/watchdog-ping-socket.js b/app/utils/watchdog-ping-socket.js index 5f98a9ba..cec509a1 100644 --- a/app/utils/watchdog-ping-socket.js +++ b/app/utils/watchdog-ping-socket.js @@ -1,19 +1,6 @@ import net from "node:net" import { APIs, BASE_API } from "../networking/APIs.js" -const onSentFailed = async ( - error, isTCP, userId, connection, - logger, sessionService, onTcpCloseCb, onWsCloseCb -) => { - logger.error(error, "[error socket send] %s", userId) - if (isTCP) { - await onTcpCloseCb(connection?.socket).catch(error => logger.error(error, "[close tcp]")) - } else { - await onWsCloseCb(connection?.socket, 10).catch(error => logger.error(error, "[close ws]")) - } - await sessionService.removeUserSession(connection?.socket, userId, connection?.deviceId).catch(error => logger.error(error, "[remove]")) -} - export const watchdogPingSocket = async (logger, sessionService, onWsCloseCb, onTcpCloseCb) => { const users = Object.keys(sessionService.activeSessions.DEVICES) @@ -27,26 +14,28 @@ export const watchdogPingSocket = async (logger, sessionService, onWsCloseCb, on } const isTCP = connection.socket instanceof net.Socket - const pingPackage = APIs[connection.socket.apiType ?? BASE_API].pingPackage() + const pingPackage = APIs[connection.socket?.apiType ?? BASE_API].pingPackage() try { if (isTCP) { - connection?.socket?.write(pingPackage, (error) => { - if (error) { - onSentFailed( - error, isTCP, userId, connection, - logger, sessionService, onWsCloseCb, onTcpCloseCb - ) - } + await new Promise((resolve, reject) => { + connection.socket?.write(pingPackage, (error) => error ? reject(error) : resolve()) }) } else { - connection?.socket?.send(pingPackage) + connection.socket?.send(pingPackage) } } catch (error) { - await onSentFailed( - error, isTCP, userId, connection, - logger, sessionService, onWsCloseCb, onTcpCloseCb - ) + logger.error(error, "[error socket send] %s", userId) + if (isTCP) { + await onTcpCloseCb(connection?.socket) + .then(() => logger.debug("[close tcp done] %s", userId)) + .catch(error => logger.error(error, "[close tcp error]")) + } else { + await onWsCloseCb(connection?.socket, 10) + .then(() => logger.debug("[close ws done] %s", userId)) + .catch(error => logger.error(error, "[close ws error]")) + } + await sessionService.removeUserSession(connection?.socket, userId, connection?.deviceId).catch(error => logger.error(error, "[remove]")) } } } From e885838e19685bcf7e9723b1aee5037daee6e70c Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Thu, 7 May 2026 04:06:25 +0300 Subject: [PATCH 14/16] add logs --- app/providers/services/session/Provider.js | 3 ++- app/providers/services/session/index.js | 7 ++++++- app/utils/watchdog-ping-socket.js | 4 +++- 3 files changed, 11 insertions(+), 3 deletions(-) diff --git a/app/providers/services/session/Provider.js b/app/providers/services/session/Provider.js index d5f5fa2e..3896583a 100644 --- a/app/providers/services/session/Provider.js +++ b/app/providers/services/session/Provider.js @@ -8,9 +8,10 @@ const name = "SessionService" class SessionServiceRegisterProvider extends RegisterProvider { register(slc) { const config = slc.use("Config") + const logger = slc.use("Logger").child("[SessionService]") const redisClient = slc.use("RedisClient") - return new SessionService(ACTIVE, config, redisClient) + return new SessionService(ACTIVE, config, logger, redisClient) } } diff --git a/app/providers/services/session/index.js b/app/providers/services/session/index.js index 1bd0b406..93c8165a 100644 --- a/app/providers/services/session/index.js +++ b/app/providers/services/session/index.js @@ -10,9 +10,10 @@ import { CONSTANTS } from "../../../constants/constants.js" */ class SessionService { - constructor(activeSessions, config, redisConnection) { + constructor(activeSessions, config, logger, redisConnection) { this.activeSessions = activeSessions this.config = config + this.logger = logger this.redisConnection = redisConnection } @@ -316,10 +317,14 @@ class SessionService { } async removeUserSession(socket, userId, deviceId) { + this.logger.debug("[removeUserSession][args]: %o", { socket: socket?.isAlive, userId, deviceId }) + userId = userId ?? this.getSessionUserId(socket) deviceId = deviceId ?? this.getDeviceId(socket, userId) const orgId = this.getSession(socket)?.organizationId + this.logger.debug("[removeUserSession][vars]: %o [session]: %o [device]: %s", { orgId, userId, deviceId }, this.getSession(socket), this.getDeviceId(socket, userId)) + const leftActiveConnections = this.getUserDevices(userId).filter(({ deviceId: activeDeviceId }) => activeDeviceId !== deviceId) if (leftActiveConnections?.length) { diff --git a/app/utils/watchdog-ping-socket.js b/app/utils/watchdog-ping-socket.js index cec509a1..fa082064 100644 --- a/app/utils/watchdog-ping-socket.js +++ b/app/utils/watchdog-ping-socket.js @@ -4,7 +4,7 @@ import { APIs, BASE_API } from "../networking/APIs.js" export const watchdogPingSocket = async (logger, sessionService, onWsCloseCb, onTcpCloseCb) => { const users = Object.keys(sessionService.activeSessions.DEVICES) - logger.debug("[run] %s", users.length) + logger.debug("[start] %s", users.length) for (const userId of users) { const connections = sessionService.activeSessions.DEVICES[userId] ?? [] @@ -39,4 +39,6 @@ export const watchdogPingSocket = async (logger, sessionService, onWsCloseCb, on } } } + + logger.debug("[finish]") } \ No newline at end of file From 36f1baa29b9d079622f17567aeca0378c215e59d Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Thu, 7 May 2026 13:35:52 +0300 Subject: [PATCH 15/16] update addUserDeviceConnection --- app/providers/services/session/index.js | 25 ++++++++++++++++++++++--- 1 file changed, 22 insertions(+), 3 deletions(-) diff --git a/app/providers/services/session/index.js b/app/providers/services/session/index.js index 93c8165a..1d14eefe 100644 --- a/app/providers/services/session/index.js +++ b/app/providers/services/session/index.js @@ -22,13 +22,17 @@ class SessionService { } addUserDeviceConnection(socket, organizationId, userId, deviceId) { - const activeConnections = this.activeSessions.DEVICES[userId] const socketsToClose = [] + let activeConnections = this.getUserDevices(userId) + const filterNotSameSocket = activeConnections.filter(connection => connection.socket !== socket) + this.activeSessions.DEVICES[userId] = filterNotSameSocket + activeConnections = this.getUserDevices(userId) + const connection = { socket: socket, deviceId, organizationId } if (activeConnections) { - const devices = activeConnections.filter((connection) => { + const otherDeviceConnections = activeConnections.filter((connection) => { if (connection.deviceId !== deviceId) { return true } else { @@ -36,7 +40,7 @@ class SessionService { return false } }) - this.activeSessions.DEVICES[userId] = [...devices, connection] + this.activeSessions.DEVICES[userId] = [...otherDeviceConnections, connection] } else { this.activeSessions.DEVICES[userId] = [connection] } @@ -325,6 +329,13 @@ class SessionService { this.logger.debug("[removeUserSession][vars]: %o [session]: %o [device]: %s", { orgId, userId, deviceId }, this.getSession(socket), this.getDeviceId(socket, userId)) + const devicesBefore = this.getUserDevices(userId).map((connection) => { + const { socket, ...connectionData } = connection + return { ...connectionData, socket: socket?.clientId } + }) + + this.logger.debug("[removeUserSession][devices][before]: %o %s", devicesBefore, devicesBefore?.length) + const leftActiveConnections = this.getUserDevices(userId).filter(({ deviceId: activeDeviceId }) => activeDeviceId !== deviceId) if (leftActiveConnections?.length) { @@ -334,6 +345,14 @@ class SessionService { } this.activeSessions.SESSIONS.delete(socket) + + const devicesAfter = this.getUserDevices(userId).map((connection) => { + const { socket, ...connectionData } = connection + return { ...connectionData, socket: socket?.clientId } + }) + + this.logger.debug("[removeUserSession][devices][after]: %o %s", devicesAfter, devicesAfter?.length) + if (!deviceId) { return } From 8cbf53f86d7fa51577944e85864bf978d951eb63 Mon Sep 17 00:00:00 2001 From: AntonSelezev Date: Tue, 12 May 2026 16:17:10 +0300 Subject: [PATCH 16/16] update package --- APIs/XMPP | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/APIs/XMPP b/APIs/XMPP index 80ab3a71..7345b506 160000 --- a/APIs/XMPP +++ b/APIs/XMPP @@ -1 +1 @@ -Subproject commit 80ab3a714e5b67ac9a0c130d2e0c651bff82c381 +Subproject commit 7345b506e738263a1ddcc46a48de73c0f9a437ef