Skip to content

Commit fb93855

Browse files
authored
fix online list request (#204)
* fix online list request * add basic cmd ws test commands * update check cmd connections * add close * add socketCloseWatchdog * add remove session * on close logs debug * add repl service * update docker-file * add ping socket * update package * fix tcp send failed * update watchdogPingSocket * add logs * update addUserDeviceConnection * update submodule
1 parent 4a9aa4f commit fb93855

12 files changed

Lines changed: 289 additions & 19 deletions

File tree

.env.example

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
NODE_ENV=development
22
STANDALONE_NODE=false
3+
WATCHDOG_PING_SOCKET_INTERVAL=60000 # 1 minute
34

45
APP_PORT=9001
56
APP_TCP_PORT=8001
@@ -11,6 +12,16 @@ MONGODB_URL=mongodb://127.0.0.1/samadb
1112

1213
REDIS_URL=redis://127.0.0.1:6379
1314

15+
# REPL
16+
# http
17+
APP_REPL_HTTP_PORT=5010
18+
APP_REPL_HTTP_ACCESS_KEY=repl-wow-key
19+
# socket
20+
APP_REPL_SOCKET_HANDLER=/tmp/net-repl.socket
21+
# file
22+
APP_REPL_FILE_IN=/tmp/pipe-repl.in
23+
APP_REPL_FILE_OUT=/tmp/pipe-repl.out
24+
1425
STORAGE_DRIVER=minio
1526

1627
# If you set STORAGE_DRIVER=s3, then fill the below envs

APIs/JSON/index.js

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,4 +35,10 @@ export default class JsonAPI extends BaseAPI {
3535
}
3636
return this.stringifyMessage(message)
3737
}
38+
39+
pingPackage() {
40+
const ping = { response: { ping: {} } }
41+
42+
return this.stringifyMessage(ping)
43+
}
3844
}

APIs/JSON/routes/packet_processor.js

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -46,18 +46,18 @@ class PacketJsonProcessor extends BasePacketProcessor {
4646
} catch (error) {
4747
logger.error(error)
4848
let errorBackMessage = null
49-
if (json.request) {
49+
if (json?.request) {
5050
errorBackMessage = {
5151
response: {
52-
id: json.request.id,
52+
id: json.request?.id,
5353
error: error.cause || error.message,
5454
},
5555
}
5656
} else {
57-
const topLevelElement = Object.keys(json)[0]
57+
const topLevelElement = json ? Object.keys(json)[0] : void 0
5858
errorBackMessage = {
5959
[topLevelElement]: {
60-
id: json[topLevelElement].id,
60+
id: json?.[topLevelElement]?.id,
6161
error: error.cause || error.message,
6262
},
6363
}

APIs/XMPP

Submodule XMPP updated from ba4c3af to 7345b50

Dockerfile

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@ COPY . .
1111
FROM node:22-slim
1212

1313
RUN apt-get update && apt-get install -y --no-install-recommends \
14-
curl \
14+
curl netcat-openbsd \
1515
&& rm -rf /var/lib/apt/lists/*
1616

1717
WORKDIR /app

app/config/default.js

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ const CONFIG = {
77
name: process.env.APP_NAME ?? "SAMA",
88
hostName: process.env.HOSTNAME,
99
isStandAloneNode: process.env.STANDALONE_NODE === CONSTANTS.ENV_TRUE,
10+
watchdogPingSocketInterval: +process.env.WATCHDOG_PING_SOCKET_INTERVAL,
1011
},
1112
logger: {
1213
logLevel: process.env.LOG_LEVEL ?? "debug",
@@ -102,6 +103,19 @@ const CONFIG = {
102103
apiKey: process.env.HTTP_ADMIN_API_KEY,
103104
},
104105
},
106+
repl: {
107+
http: {
108+
port: +(process.env.APP_REPL_HTTP_PORT ?? 5010),
109+
accessKey: process.env.APP_REPL_HTTP_ACCESS_KEY,
110+
},
111+
socket: {
112+
handler: process.env.APP_REPL_SOCKET_HANDLER,
113+
},
114+
file: {
115+
in: process.env.APP_REPL_FILE_IN,
116+
out: process.env.APP_REPL_FILE_OUT,
117+
},
118+
},
105119
conversation: {
106120
disableChannelsLogic: process.env.CONVERSATION_DISABLE_CHANNELS_LOGIC === CONSTANTS.ENV_TRUE,
107121
isEventsEnabled: process.env.CONVERSATION_NOTIFICATIONS_ENABLED === CONSTANTS.ENV_TRUE,

app/lib/repl-tools.js

Lines changed: 144 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,144 @@
1+
import fs from "node:fs"
2+
import fsPromises from "node:fs/promises"
3+
import vm from "node:vm"
4+
import { exec } from "node:child_process"
5+
import repl from "node:repl"
6+
import net from "node:net"
7+
import http from "node:http"
8+
9+
import { CONSTANTS } from "../constants/constants.js"
10+
11+
// example: curl --no-progress-meter -sSNT. -H "api-key: ****" localhost:5010
12+
13+
const httpReplService = async (replOptions, httpOptions) => {
14+
const { ctx } = replOptions
15+
const { accessKey, port } = httpOptions
16+
17+
const server = http.createServer((req, res) => {
18+
if (req.headers[CONSTANTS.HTTP_REPL_ACCESS_KEY_HEADER] !== accessKey) {
19+
return res.end("Invalid access-key")
20+
}
21+
22+
res.setHeader("content-type", "multipart/octet-stream")
23+
24+
const replService = repl.start({
25+
prompt: "curl repl> ",
26+
input: req,
27+
output: res,
28+
terminal: false,
29+
useColors: true,
30+
useGlobal: false,
31+
})
32+
33+
const context = vm.createContext({ ...ctx })
34+
35+
replService.context = context
36+
37+
req.on("error", (error) => replService.close())
38+
req.on("end", () => replService.close())
39+
40+
res.on("error", (error) => replService.close())
41+
res.on("end", () => replService.close())
42+
43+
replService.on("close", () => !res.closed && res.end("REPL closed"))
44+
replService.on("error", (error) => {
45+
replService.close()
46+
!res.closed && res.end("REPL closed")
47+
!req.closed && req.destroy()
48+
})
49+
})
50+
51+
return new Promise((resolve, reject) => {
52+
server.listen(port, () => resolve(port))
53+
})
54+
}
55+
56+
// example write: cat > ./pipe-repl.in
57+
// example read: tail -f ./pipe-repl.out
58+
59+
const fileReplService = async (replOptions, fileOptions) => {
60+
const { ctx } = replOptions
61+
const { fileIn, fileOut } = fileOptions
62+
63+
await fsPromises.rm(fileIn, { force: true }).catch((error) => {})
64+
await fsPromises.rm(fileOut, { force: true }).catch((error) => {})
65+
66+
await new Promise((resolve, reject) => {
67+
exec(`mkfifo ${fileIn}`, (error, out, outErr) => {
68+
if (error) {
69+
return reject(error)
70+
}
71+
resolve()
72+
})
73+
})
74+
75+
const cmdInputPipe = fs.createReadStream(fileIn, { encoding: "utf8" })
76+
const cmdOutPipe = fs.createWriteStream(fileOut, { encoding: "utf8" })
77+
78+
const replService = repl.start({
79+
prompt: "file repl> ",
80+
input: cmdInputPipe,
81+
output: cmdOutPipe,
82+
terminal: false,
83+
useColors: true,
84+
useGlobal: false,
85+
})
86+
87+
const context = vm.createContext({ ...ctx })
88+
89+
replService.context = context
90+
91+
replService.on("error", (error) => {
92+
replService.close()
93+
cmdInputPipe.close()
94+
})
95+
96+
cmdInputPipe.on("end", () => {
97+
replService.close()
98+
fileReplService(replOptions, fileOptions)
99+
})
100+
}
101+
102+
// example: nc -U ./net-repl.socket
103+
104+
const netReplService = async (replOptions, fileOptions) => {
105+
const { ctx } = replOptions
106+
const { socketHandler } = fileOptions
107+
108+
await fsPromises.rm(socketHandler, { force: true }).catch((error) => {})
109+
110+
const server = net.createServer((socket) => {
111+
const replService = repl.start({
112+
prompt: "socket repl> ",
113+
input: socket,
114+
output: socket,
115+
terminal: false,
116+
useColors: true,
117+
useGlobal: false,
118+
})
119+
120+
const context = vm.createContext({ ...ctx })
121+
122+
replService.context = context
123+
124+
socket.on("end", () => replService.close())
125+
socket.on("error", (error) => replService.close())
126+
replService.on("close", () => !socket.closed && socket.end("REPL closed"))
127+
})
128+
129+
return new Promise((resolve, reject) => {
130+
server.listen(socketHandler, () => resolve(socketHandler))
131+
})
132+
}
133+
134+
export const startReplServices = async (replOptions, httpOptions, netOptions, fileOptions) => {
135+
if (httpOptions.accessKey) {
136+
await httpReplService(replOptions, httpOptions)
137+
}
138+
if (netOptions.socketHandler) {
139+
await netReplService(replOptions, netOptions)
140+
}
141+
if (fileOptions.fileIn && fileOptions.fileOut) {
142+
await fileReplService(replOptions, fileOptions)
143+
}
144+
}

app/networking/protocol_processors/base.js

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -200,7 +200,7 @@ class BaseProtocolProcessor {
200200
}
201201

202202
async onClose(socket, code) {
203-
logger.trace("[Close] IP: %s CLIENT_ID: %s CODE: %s", this.socketAddress(socket), socket.clientId, code)
203+
logger.debug("[Close] IP: %s CLIENT_ID: %s CODE: %s", this.socketAddress(socket), socket.clientId, code)
204204

205205
socket.isAlive = false
206206

@@ -214,7 +214,7 @@ class BaseProtocolProcessor {
214214
async updateLastUserLastActivityOnClose(socket) {
215215
const { organizationId, userId } = this.sessionService.getSession(socket) ?? {}
216216

217-
logger.trace("[UPDATE_LAST_ACTIVITY][CLOSE] OrgId: %s UserId: %s", organizationId, userId)
217+
logger.debug("[UPDATE_LAST_ACTIVITY][CLOSE] OrgId: %s UserId: %s", organizationId, userId)
218218

219219
if (!userId) {
220220
return

app/providers/services/session/Provider.js

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,9 +8,10 @@ const name = "SessionService"
88
class SessionServiceRegisterProvider extends RegisterProvider {
99
register(slc) {
1010
const config = slc.use("Config")
11+
const logger = slc.use("Logger").child("[SessionService]")
1112
const redisClient = slc.use("RedisClient")
1213

13-
return new SessionService(ACTIVE, config, redisClient)
14+
return new SessionService(ACTIVE, config, logger, redisClient)
1415
}
1516
}
1617

app/providers/services/session/index.js

Lines changed: 36 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -10,9 +10,10 @@ import { CONSTANTS } from "../../../constants/constants.js"
1010
*/
1111

1212
class SessionService {
13-
constructor(activeSessions, config, redisConnection) {
13+
constructor(activeSessions, config, logger, redisConnection) {
1414
this.activeSessions = activeSessions
1515
this.config = config
16+
this.logger = logger
1617
this.redisConnection = redisConnection
1718
}
1819

@@ -21,21 +22,25 @@ class SessionService {
2122
}
2223

2324
addUserDeviceConnection(socket, organizationId, userId, deviceId) {
24-
const activeConnections = this.activeSessions.DEVICES[userId]
2525
const socketsToClose = []
2626

27+
let activeConnections = this.getUserDevices(userId)
28+
const filterNotSameSocket = activeConnections.filter(connection => connection.socket !== socket)
29+
this.activeSessions.DEVICES[userId] = filterNotSameSocket
30+
activeConnections = this.getUserDevices(userId)
31+
2732
const connection = { socket: socket, deviceId, organizationId }
2833

2934
if (activeConnections) {
30-
const devices = activeConnections.filter((connection) => {
35+
const otherDeviceConnections = activeConnections.filter((connection) => {
3136
if (connection.deviceId !== deviceId) {
3237
return true
3338
} else {
3439
socketsToClose.push(connection.socket)
3540
return false
3641
}
3742
})
38-
this.activeSessions.DEVICES[userId] = [...devices, connection]
43+
this.activeSessions.DEVICES[userId] = [...otherDeviceConnections, connection]
3944
} else {
4045
this.activeSessions.DEVICES[userId] = [connection]
4146
}
@@ -102,9 +107,7 @@ class SessionService {
102107

103108
async listUserDevice(organizationId, userId) {
104109
if (this.config.get("app.isStandAloneNode")) {
105-
return this.getUserDevices(userId)
106-
.map((connection) => connection?.deviceId)
107-
.filter((deviceId) => deviceId !== CONSTANTS.HTTP_DEVICE_ID)
110+
return this.listUserDeviceLocal(userId)
108111
}
109112

110113
const userKey = this.#usersSetCacheKey(organizationId, userId)
@@ -113,6 +116,12 @@ class SessionService {
113116
return deviceIds ?? []
114117
}
115118

119+
listUserDeviceLocal(userId) {
120+
return this.getUserDevices(userId)
121+
.map((connection) => connection?.deviceId)
122+
.filter((deviceId) => deviceId !== CONSTANTS.HTTP_DEVICE_ID)
123+
}
124+
116125
async deleteUserDevices(organizationId, userId) {
117126
const userKey = this.#usersSetCacheKey(organizationId, userId)
118127

@@ -312,10 +321,21 @@ class SessionService {
312321
}
313322

314323
async removeUserSession(socket, userId, deviceId) {
324+
this.logger.debug("[removeUserSession][args]: %o", { socket: socket?.isAlive, userId, deviceId })
325+
315326
userId = userId ?? this.getSessionUserId(socket)
316327
deviceId = deviceId ?? this.getDeviceId(socket, userId)
317328
const orgId = this.getSession(socket)?.organizationId
318329

330+
this.logger.debug("[removeUserSession][vars]: %o [session]: %o [device]: %s", { orgId, userId, deviceId }, this.getSession(socket), this.getDeviceId(socket, userId))
331+
332+
const devicesBefore = this.getUserDevices(userId).map((connection) => {
333+
const { socket, ...connectionData } = connection
334+
return { ...connectionData, socket: socket?.clientId }
335+
})
336+
337+
this.logger.debug("[removeUserSession][devices][before]: %o %s", devicesBefore, devicesBefore?.length)
338+
319339
const leftActiveConnections = this.getUserDevices(userId).filter(({ deviceId: activeDeviceId }) => activeDeviceId !== deviceId)
320340

321341
if (leftActiveConnections?.length) {
@@ -325,6 +345,14 @@ class SessionService {
325345
}
326346
this.activeSessions.SESSIONS.delete(socket)
327347

348+
349+
const devicesAfter = this.getUserDevices(userId).map((connection) => {
350+
const { socket, ...connectionData } = connection
351+
return { ...connectionData, socket: socket?.clientId }
352+
})
353+
354+
this.logger.debug("[removeUserSession][devices][after]: %o %s", devicesAfter, devicesAfter?.length)
355+
328356
if (!deviceId) {
329357
return
330358
}
@@ -393,7 +421,7 @@ class SessionService {
393421
(session) =>
394422
session?.organizationId === organizationId &&
395423
session?.extraParams[CONSTANTS.SESSION_DEVICE_ID_KEY] !== CONSTANTS.HTTP_DEVICE_ID &&
396-
session?.userId
424+
session?.userId && this.listUserDeviceLocal(session?.userId)?.length
397425
)
398426
.map((session) => session.userId)
399427
.sort((userIdA, userIdB) => userIdA - userIdB)

0 commit comments

Comments
 (0)