diff --git a/.env.example b/.env.example index e51f056b..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 @@ -11,6 +12,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/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/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/APIs/XMPP b/APIs/XMPP index ba4c3af7..7345b506 160000 --- a/APIs/XMPP +++ b/APIs/XMPP @@ -1 +1 @@ -Subproject commit ba4c3af760d78fb9da77e64949e6186b99878f0f +Subproject commit 7345b506e738263a1ddcc46a48de73c0f9a437ef 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 diff --git a/app/config/default.js b/app/config/default.js index 5ea61b78..c73ac88e 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, + watchdogPingSocketInterval: +process.env.WATCHDOG_PING_SOCKET_INTERVAL, }, logger: { logLevel: process.env.LOG_LEVEL ?? "debug", @@ -102,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/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 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 55f0ce29..1d14eefe 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 } @@ -21,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 { @@ -35,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] } @@ -102,9 +107,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 +116,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) @@ -312,10 +321,21 @@ 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 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) { @@ -325,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 } @@ -393,7 +421,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) diff --git a/app/utils/watchdog-ping-socket.js b/app/utils/watchdog-ping-socket.js new file mode 100644 index 00000000..fa082064 --- /dev/null +++ b/app/utils/watchdog-ping-socket.js @@ -0,0 +1,44 @@ +import net from "node:net" +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("[start] %s", users.length) + + for (const userId of users) { + const connections = sessionService.activeSessions.DEVICES[userId] ?? [] + for (const connection of connections) { + if (!connection?.socket) { + continue + } + + const isTCP = connection.socket instanceof net.Socket + const pingPackage = APIs[connection.socket?.apiType ?? BASE_API].pingPackage() + + try { + if (isTCP) { + await new Promise((resolve, reject) => { + connection.socket?.write(pingPackage, (error) => error ? reject(error) : resolve()) + }) + } else { + connection.socket?.send(pingPackage) + } + } catch (error) { + 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]")) + } + } + } + + logger.debug("[finish]") +} \ No newline at end of file diff --git a/index.js b/index.js index fee58008..ca0b0de2 100644 --- a/index.js +++ b/index.js @@ -20,10 +20,12 @@ 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" import { buildWsEndpoint } from "./app/utils/build_ws_endpoint.js" +import { watchdogPingSocket } from "./app/utils/watchdog-ping-socket.js" if (config.get("app.env") === CONSTANTS.ENVS.PROD) { process.on("unhandledRejection", (reason, promise) => { @@ -190,9 +192,29 @@ 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) } -// https://dev.to/mattkrick/replacing-express-with-uwebsockets-48ph +if (config.get("app.watchdogPingSocketInterval")) { + const socketCloseWatchdogLogger = logger.child("[SocketClosedWatchDog]") + setInterval(() => { + watchdogPingSocket( + socketCloseWatchdogLogger, + sessionService, + (socket, code) => wsProtocolImp.onClose(socket, code), + (socket) => tcpProtocolImp.onClose(socket) + ) + }, config.get("app.watchdogPingSocketInterval")) +} + +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") } +)