feat: add websocket scheme
This commit is contained in:
parent
3015f13e0d
commit
8904fa3fb8
9 changed files with 178 additions and 21 deletions
12
.env.example
12
.env.example
|
|
@ -1,9 +1,9 @@
|
||||||
## Application ##
|
## Application ##
|
||||||
PORT=3000 # OPTIONAL, DEFAULT 3000
|
PORT=3000 # OPTIONAL, DEFAULT 3000
|
||||||
API_KEY=your_global_api_key_here # OPTIONAL, DEFAULT EMPTY
|
API_KEY=SET_YOUR_API_KEY_HERE # OPTIONAL, IF SET, ALL REQUESTS MUST INCLUDE THIS IN THE 'x-api-key' HEADER
|
||||||
BASE_WEBHOOK_URL=http://localhost:3000/localCallbackExample # MANDATORY
|
BASE_WEBHOOK_URL=http://localhost:3000/localCallbackExample # MANDATORY
|
||||||
ENABLE_LOCAL_CALLBACK_EXAMPLE=TRUE # OPTIONAL, DISABLE FOR PRODUCTION
|
ENABLE_LOCAL_CALLBACK_EXAMPLE=TRUE # OPTIONAL, DISABLE FOR PRODUCTION
|
||||||
RATE_LIMIT_MAX=1000 # OPTIONAL, THE MAXIUM NUMBER OF CONNECTIONS TO ALLOW PER TIME FRAME
|
RATE_LIMIT_MAX=1000 # OPTIONAL, THE MAXIMUM NUMBER OF CONNECTIONS TO ALLOW PER TIME FRAME
|
||||||
RATE_LIMIT_WINDOW_MS=1000 # OPTIONAL, TIME FRAME FOR WHICH REQUESTS ARE CHECKED IN MS
|
RATE_LIMIT_WINDOW_MS=1000 # OPTIONAL, TIME FRAME FOR WHICH REQUESTS ARE CHECKED IN MS
|
||||||
|
|
||||||
## Client ##
|
## Client ##
|
||||||
|
|
@ -12,13 +12,15 @@ SET_MESSAGES_AS_SEEN=TRUE # WILL MARK THE MESSAGES AS READ AUTOMATICALLY
|
||||||
# ALL CALLBACKS: auth_failure|authenticated|call|change_state|disconnected|group_join|group_leave|group_update|loading_screen|media_uploaded|message|message_ack|message_create|message_reaction|message_revoke_everyone|qr|ready|contact_changed|unread_count|message_edit|message_ciphertext
|
# ALL CALLBACKS: auth_failure|authenticated|call|change_state|disconnected|group_join|group_leave|group_update|loading_screen|media_uploaded|message|message_ack|message_create|message_reaction|message_revoke_everyone|qr|ready|contact_changed|unread_count|message_edit|message_ciphertext
|
||||||
DISABLED_CALLBACKS=message_ack|message_reaction|unread_count|message_edit|message_ciphertext # PREVENT SENDING CERTAIN TYPES OF CALLBACKS BACK TO THE WEBHOOK
|
DISABLED_CALLBACKS=message_ack|message_reaction|unread_count|message_edit|message_ciphertext # PREVENT SENDING CERTAIN TYPES OF CALLBACKS BACK TO THE WEBHOOK
|
||||||
WEB_VERSION='2.2328.5' # OPTIONAL, THE VERSION OF WHATSAPP WEB TO USE
|
WEB_VERSION='2.2328.5' # OPTIONAL, THE VERSION OF WHATSAPP WEB TO USE
|
||||||
WEB_VERSION_CACHE_TYPE=none # OPTIONAL, DETERMINTES WHERE TO GET THE WHATSAPP WEB VERSION(local, remote or none), DEFAULT 'none'
|
WEB_VERSION_CACHE_TYPE=none # OPTIONAL, DETERMINES WHERE TO GET THE WHATSAPP WEB VERSION(local, remote or none), DEFAULT 'none'
|
||||||
RECOVER_SESSIONS=TRUE # OPTIONAL, SHOULD WE RECOVER THE SESSION IN CASE OF PAGE FAILURES
|
RECOVER_SESSIONS=TRUE # OPTIONAL, SHOULD WE RECOVER THE SESSION IN CASE OF PAGE FAILURES
|
||||||
CHROME_BIN= # OPTIONAL, PATH TO CHROME BINARY
|
CHROME_BIN= # OPTIONAL, PATH TO CHROME BINARY
|
||||||
HEADLESS=TRUE # OPTIONAL, RUN CHROME IN HEADLESS MODE
|
HEADLESS=TRUE # OPTIONAL, RUN CHROME IN HEADLESS MODE
|
||||||
RELEASE_BROWSER_LOCK=TRUE # OPTIONAL, RELEASE THE BROWSER LOCK ON SESSION INITIALIZATION
|
RELEASE_BROWSER_LOCK=TRUE # OPTIONAL, RELEASE THE BROWSER LOCK ON SESSION INITIALIZATION
|
||||||
|
LOGLEVEL=info # OPTIONAL, SET THE LOG LEVEL
|
||||||
|
ENABLE_WEBHOOK=TRUE # OPTIONAL, ENABLE WEBHOOK FOR REALTIME UPDATES(TRUE BY DEFAULT)
|
||||||
|
ENABLE_WEBSOCKET=FALSE # OPTIONAL, ENABLE WEBSOCKET FOR REALTIME UPDATES(FALSE BY DEFAULT)
|
||||||
|
|
||||||
## Session File Storage ##
|
## Session File Storage ##
|
||||||
SESSIONS_PATH=./sessions # OPTIONAL
|
SESSIONS_PATH=./sessions # OPTIONAL
|
||||||
|
ENABLE_SWAGGER_ENDPOINT=TRUE # OPTIONAL, ENABLE SWAGGER ENDPOINT FOR API DOCUMENTATION
|
||||||
ENABLE_SWAGGER_ENDPOINT=TRUE # OPTIONAL
|
|
||||||
13
README.md
13
README.md
|
|
@ -149,10 +149,23 @@ For example, if you have the sessionId defined as `DEMO`, the environment variab
|
||||||
|
|
||||||
By setting the `DISABLED_CALLBACKS` environment variable you can specify what events you are **not** willing to receive on your webhook.
|
By setting the `DISABLED_CALLBACKS` environment variable you can specify what events you are **not** willing to receive on your webhook.
|
||||||
|
|
||||||
|
By setting the `ENABLE_WEBHOOK` environment to `FALSE` you can disable webhook dispatching. This will help you if you want to switch to websocket method(see below).
|
||||||
|
|
||||||
### Scanning QR code
|
### Scanning QR code
|
||||||
|
|
||||||
In order to validate a new WhatsApp Web instance you need to scan the QR code using your mobile phone. Official documentation can be found at (https://faq.whatsapp.com/1079327266110265/?cms_platform=android) page. The service itself delivers the QR code content as a webhook event or you can use the REST endpoints (`/session/qr/:sessionId` or `/session/qr/:sessionId/image` to get the QR code as a png image).
|
In order to validate a new WhatsApp Web instance you need to scan the QR code using your mobile phone. Official documentation can be found at (https://faq.whatsapp.com/1079327266110265/?cms_platform=android) page. The service itself delivers the QR code content as a webhook event or you can use the REST endpoints (`/session/qr/:sessionId` or `/session/qr/:sessionId/image` to get the QR code as a png image).
|
||||||
|
|
||||||
|
### WebSocket mode
|
||||||
|
The service can dispatch realtime events through websocket connection. By default, the websocket is not activated, so you need manually set the `ENABLE_WEBSOCKET` environment variable to activate it. The server activates a new websocket instance per each active session. The websocket path is `/ws/:sessionId`, where sessionId is your configured session name. The websocket supports ping/pong scheme to keep the socket running.
|
||||||
|
The below example shows how to receive the events for **test** session.
|
||||||
|
```
|
||||||
|
const ws = new WebSocket('ws://127.0.0.1:3000/ws/test');
|
||||||
|
|
||||||
|
ws.on('message', (data) => {
|
||||||
|
// consume the events
|
||||||
|
});
|
||||||
|
```
|
||||||
|
|
||||||
## Deploy to Production
|
## Deploy to Production
|
||||||
|
|
||||||
- Load the docker image in docker-compose, or your Kubernetes environment
|
- Load the docker image in docker-compose, or your Kubernetes environment
|
||||||
|
|
|
||||||
32
package-lock.json
generated
32
package-lock.json
generated
|
|
@ -17,7 +17,8 @@
|
||||||
"qr-image": "^3.2.0",
|
"qr-image": "^3.2.0",
|
||||||
"qrcode-terminal": "^0.12.0",
|
"qrcode-terminal": "^0.12.0",
|
||||||
"swagger-ui-express": "^5.0.1",
|
"swagger-ui-express": "^5.0.1",
|
||||||
"whatsapp-web.js": "^1.26.1-alpha.3"
|
"whatsapp-web.js": "^1.26.1-alpha.3",
|
||||||
|
"ws": "^8.18.0"
|
||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"eslint": "^8.38.0",
|
"eslint": "^8.38.0",
|
||||||
|
|
@ -6656,6 +6657,27 @@
|
||||||
"integrity": "sha512-sGkPx+VjMtmA6MX27oA4FBFELFCZZ4S4XqeGOXCv68tT+jb3vk/RyaKWP0PTKyWtmLSM0b+adUTEvbs1PEaH2w==",
|
"integrity": "sha512-sGkPx+VjMtmA6MX27oA4FBFELFCZZ4S4XqeGOXCv68tT+jb3vk/RyaKWP0PTKyWtmLSM0b+adUTEvbs1PEaH2w==",
|
||||||
"license": "MIT"
|
"license": "MIT"
|
||||||
},
|
},
|
||||||
|
"node_modules/puppeteer-core/node_modules/ws": {
|
||||||
|
"version": "8.9.0",
|
||||||
|
"resolved": "https://registry.npmjs.org/ws/-/ws-8.9.0.tgz",
|
||||||
|
"integrity": "sha512-Ja7nszREasGaYUYCI2k4lCKIRTt+y7XuqVoHR44YpI49TtryyqbqvDMn5eqfW7e6HzTukDRIsXqzVHScqRcafg==",
|
||||||
|
"license": "MIT",
|
||||||
|
"engines": {
|
||||||
|
"node": ">=10.0.0"
|
||||||
|
},
|
||||||
|
"peerDependencies": {
|
||||||
|
"bufferutil": "^4.0.1",
|
||||||
|
"utf-8-validate": "^5.0.2"
|
||||||
|
},
|
||||||
|
"peerDependenciesMeta": {
|
||||||
|
"bufferutil": {
|
||||||
|
"optional": true
|
||||||
|
},
|
||||||
|
"utf-8-validate": {
|
||||||
|
"optional": true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
"node_modules/pure-rand": {
|
"node_modules/pure-rand": {
|
||||||
"version": "6.1.0",
|
"version": "6.1.0",
|
||||||
"resolved": "https://registry.npmjs.org/pure-rand/-/pure-rand-6.1.0.tgz",
|
"resolved": "https://registry.npmjs.org/pure-rand/-/pure-rand-6.1.0.tgz",
|
||||||
|
|
@ -8274,16 +8296,16 @@
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"node_modules/ws": {
|
"node_modules/ws": {
|
||||||
"version": "8.9.0",
|
"version": "8.18.0",
|
||||||
"resolved": "https://registry.npmjs.org/ws/-/ws-8.9.0.tgz",
|
"resolved": "https://registry.npmjs.org/ws/-/ws-8.18.0.tgz",
|
||||||
"integrity": "sha512-Ja7nszREasGaYUYCI2k4lCKIRTt+y7XuqVoHR44YpI49TtryyqbqvDMn5eqfW7e6HzTukDRIsXqzVHScqRcafg==",
|
"integrity": "sha512-8VbfWfHLbbwu3+N6OKsOMpBdT4kXPDDB9cJk2bJ6mh9ucxdlnNvH1e+roYkKmN9Nxw2yjz7VzeO9oOz2zJ04Pw==",
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
"engines": {
|
"engines": {
|
||||||
"node": ">=10.0.0"
|
"node": ">=10.0.0"
|
||||||
},
|
},
|
||||||
"peerDependencies": {
|
"peerDependencies": {
|
||||||
"bufferutil": "^4.0.1",
|
"bufferutil": "^4.0.1",
|
||||||
"utf-8-validate": "^5.0.2"
|
"utf-8-validate": ">=5.0.2"
|
||||||
},
|
},
|
||||||
"peerDependenciesMeta": {
|
"peerDependenciesMeta": {
|
||||||
"bufferutil": {
|
"bufferutil": {
|
||||||
|
|
|
||||||
|
|
@ -17,7 +17,8 @@
|
||||||
"qr-image": "^3.2.0",
|
"qr-image": "^3.2.0",
|
||||||
"qrcode-terminal": "^0.12.0",
|
"qrcode-terminal": "^0.12.0",
|
||||||
"swagger-ui-express": "^5.0.1",
|
"swagger-ui-express": "^5.0.1",
|
||||||
"whatsapp-web.js": "^1.26.1-alpha.3"
|
"whatsapp-web.js": "^1.26.1-alpha.3",
|
||||||
|
"ws": "^8.18.0"
|
||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"eslint": "^8.38.0",
|
"eslint": "^8.38.0",
|
||||||
|
|
|
||||||
16
server.js
16
server.js
|
|
@ -1,21 +1,29 @@
|
||||||
const app = require('./src/app')
|
const app = require('./src/app')
|
||||||
const { baseWebhookURL } = require('./src/config')
|
const { baseWebhookURL, enableWebHook, enableWebSocket } = require('./src/config')
|
||||||
const { logger } = require('./src/logger')
|
const { logger } = require('./src/logger')
|
||||||
|
const { handleUpgrade } = require('./src/websocket')
|
||||||
|
|
||||||
require('dotenv').config()
|
require('dotenv').config()
|
||||||
|
|
||||||
// Start the server
|
// Start the server
|
||||||
const port = process.env.PORT || 3000
|
const port = process.env.PORT || 3000
|
||||||
|
|
||||||
// Check if BASE_WEBHOOK_URL environment variable is available
|
// Check if BASE_WEBHOOK_URL environment variable is available when WebHook is enabled
|
||||||
if (!baseWebhookURL) {
|
if (!baseWebhookURL && enableWebHook) {
|
||||||
logger.error('BASE_WEBHOOK_URL environment variable is not set. Exiting...')
|
logger.error('BASE_WEBHOOK_URL environment variable is not set. Exiting...')
|
||||||
process.exit(1) // Terminate the application with an error code
|
process.exit(1) // Terminate the application with an error code
|
||||||
}
|
}
|
||||||
|
|
||||||
app.listen(port, () => {
|
const server = app.listen(port, () => {
|
||||||
logger.info(`Server running on port ${port}`)
|
logger.info(`Server running on port ${port}`)
|
||||||
})
|
})
|
||||||
|
|
||||||
|
if (enableWebSocket) {
|
||||||
|
server.on('upgrade', (request, socket, head) => {
|
||||||
|
handleUpgrade(request, socket, head)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
// puppeteer uses subscriptions to SIGINT, SIGTERM, and SIGHUP to know when to close browser instances
|
// puppeteer uses subscriptions to SIGINT, SIGTERM, and SIGHUP to know when to close browser instances
|
||||||
// this disables the warnings when you starts more than 10 browser instances
|
// this disables the warnings when you starts more than 10 browser instances
|
||||||
process.setMaxListeners(0)
|
process.setMaxListeners(0)
|
||||||
|
|
|
||||||
|
|
@ -19,6 +19,8 @@ const chromeBin = process.env.CHROME_BIN || null
|
||||||
const headless = process.env.HEADLESS ? (process.env.HEADLESS).toLowerCase() === 'true' : true
|
const headless = process.env.HEADLESS ? (process.env.HEADLESS).toLowerCase() === 'true' : true
|
||||||
const releaseBrowserLock = process.env.RELEASE_BROWSER_LOCK ? (process.env.RELEASE_BROWSER_LOCK).toLowerCase() === 'true' : true
|
const releaseBrowserLock = process.env.RELEASE_BROWSER_LOCK ? (process.env.RELEASE_BROWSER_LOCK).toLowerCase() === 'true' : true
|
||||||
const logLevel = process.env.LOGLEVEL || 'info'
|
const logLevel = process.env.LOGLEVEL || 'info'
|
||||||
|
const enableWebHook = process.env.ENABLE_WEBHOOK ? (process.env.ENABLE_WEBHOOK).toLowerCase() === 'true' : true
|
||||||
|
const enableWebSocket = process.env.ENABLE_WEBSOCKET ? (process.env.ENABLE_WEBSOCKET).toLowerCase() === 'true' : false
|
||||||
|
|
||||||
module.exports = {
|
module.exports = {
|
||||||
sessionFolderPath,
|
sessionFolderPath,
|
||||||
|
|
@ -37,5 +39,7 @@ module.exports = {
|
||||||
chromeBin,
|
chromeBin,
|
||||||
headless,
|
headless,
|
||||||
releaseBrowserLock,
|
releaseBrowserLock,
|
||||||
logLevel
|
logLevel,
|
||||||
|
enableWebHook,
|
||||||
|
enableWebSocket
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -5,6 +5,7 @@ const sessions = new Map()
|
||||||
const { baseWebhookURL, sessionFolderPath, maxAttachmentSize, setMessagesAsSeen, webVersion, webVersionCacheType, recoverSessions, chromeBin, headless, releaseBrowserLock } = require('./config')
|
const { baseWebhookURL, sessionFolderPath, maxAttachmentSize, setMessagesAsSeen, webVersion, webVersionCacheType, recoverSessions, chromeBin, headless, releaseBrowserLock } = require('./config')
|
||||||
const { triggerWebhook, waitForNestedObject, checkIfEventisEnabled, sendMessageSeenStatus } = require('./utils')
|
const { triggerWebhook, waitForNestedObject, checkIfEventisEnabled, sendMessageSeenStatus } = require('./utils')
|
||||||
const { logger } = require('./logger')
|
const { logger } = require('./logger')
|
||||||
|
const { initWebSocketServer, terminateWebSocketServer, triggerWebSocket } = require('./websocket')
|
||||||
|
|
||||||
// Function to validate if the session is ready
|
// Function to validate if the session is ready
|
||||||
const validateSession = async (sessionId) => {
|
const validateSession = async (sessionId) => {
|
||||||
|
|
@ -143,6 +144,7 @@ const setupSession = async (sessionId) => {
|
||||||
throw error
|
throw error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
initWebSocketServer(sessionId)
|
||||||
initializeEvents(client, sessionId)
|
initializeEvents(client, sessionId)
|
||||||
|
|
||||||
// Save the session to the Map
|
// Save the session to the Map
|
||||||
|
|
@ -181,6 +183,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('auth_failure', (msg) => {
|
client.on('auth_failure', (msg) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'status', { msg })
|
triggerWebhook(sessionWebhook, sessionId, 'status', { msg })
|
||||||
|
triggerWebSocket(sessionId, 'status', { msg })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -188,6 +191,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('authenticated', () => {
|
client.on('authenticated', () => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'authenticated')
|
triggerWebhook(sessionWebhook, sessionId, 'authenticated')
|
||||||
|
triggerWebSocket(sessionId, 'authenticated')
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -195,6 +199,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('call', async (call) => {
|
client.on('call', async (call) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'call', { call })
|
triggerWebhook(sessionWebhook, sessionId, 'call', { call })
|
||||||
|
triggerWebSocket(sessionId, 'call', { call })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -202,6 +207,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('change_state', state => {
|
client.on('change_state', state => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'change_state', { state })
|
triggerWebhook(sessionWebhook, sessionId, 'change_state', { state })
|
||||||
|
triggerWebSocket(sessionId, 'change_state', { state })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -209,6 +215,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('disconnected', (reason) => {
|
client.on('disconnected', (reason) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'disconnected', { reason })
|
triggerWebhook(sessionWebhook, sessionId, 'disconnected', { reason })
|
||||||
|
triggerWebSocket(sessionId, 'disconnected', { reason })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -216,6 +223,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('group_join', (notification) => {
|
client.on('group_join', (notification) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'group_join', { notification })
|
triggerWebhook(sessionWebhook, sessionId, 'group_join', { notification })
|
||||||
|
triggerWebSocket(sessionId, 'group_join', { notification })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -223,6 +231,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('group_leave', (notification) => {
|
client.on('group_leave', (notification) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'group_leave', { notification })
|
triggerWebhook(sessionWebhook, sessionId, 'group_leave', { notification })
|
||||||
|
triggerWebSocket(sessionId, 'group_leave', { notification })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -230,6 +239,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('group_admin_changed', (notification) => {
|
client.on('group_admin_changed', (notification) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'group_admin_changed', { notification })
|
triggerWebhook(sessionWebhook, sessionId, 'group_admin_changed', { notification })
|
||||||
|
triggerWebSocket(sessionId, 'group_admin_changed', { notification })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -237,6 +247,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('group_membership_request', (notification) => {
|
client.on('group_membership_request', (notification) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'group_membership_request', { notification })
|
triggerWebhook(sessionWebhook, sessionId, 'group_membership_request', { notification })
|
||||||
|
triggerWebSocket(sessionId, 'group_membership_request', { notification })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -244,6 +255,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('group_update', (notification) => {
|
client.on('group_update', (notification) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'group_update', { notification })
|
triggerWebhook(sessionWebhook, sessionId, 'group_update', { notification })
|
||||||
|
triggerWebSocket(sessionId, 'group_update', { notification })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -251,6 +263,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('loading_screen', (percent, message) => {
|
client.on('loading_screen', (percent, message) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'loading_screen', { percent, message })
|
triggerWebhook(sessionWebhook, sessionId, 'loading_screen', { percent, message })
|
||||||
|
triggerWebSocket(sessionId, 'loading_screen', { percent, message })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -258,6 +271,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('media_uploaded', (message) => {
|
client.on('media_uploaded', (message) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'media_uploaded', { message })
|
triggerWebhook(sessionWebhook, sessionId, 'media_uploaded', { message })
|
||||||
|
triggerWebSocket(sessionId, 'media_uploaded', { message })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -265,11 +279,13 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('message', async (message) => {
|
client.on('message', async (message) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'message', { message })
|
triggerWebhook(sessionWebhook, sessionId, 'message', { message })
|
||||||
|
triggerWebSocket(sessionId, 'message', { message })
|
||||||
if (message.hasMedia && message._data?.size < maxAttachmentSize) {
|
if (message.hasMedia && message._data?.size < maxAttachmentSize) {
|
||||||
// custom service event
|
// custom service event
|
||||||
checkIfEventisEnabled('media').then(_ => {
|
checkIfEventisEnabled('media').then(_ => {
|
||||||
message.downloadMedia().then(messageMedia => {
|
message.downloadMedia().then(messageMedia => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'media', { messageMedia, message })
|
triggerWebhook(sessionWebhook, sessionId, 'media', { messageMedia, message })
|
||||||
|
triggerWebSocket(sessionId, 'media', { messageMedia, message })
|
||||||
}).catch(error => {
|
}).catch(error => {
|
||||||
logger.error({ sessionId, err: error }, 'Failed to download media')
|
logger.error({ sessionId, err: error }, 'Failed to download media')
|
||||||
})
|
})
|
||||||
|
|
@ -285,6 +301,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('message_ack', async (message, ack) => {
|
client.on('message_ack', async (message, ack) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'message_ack', { message, ack })
|
triggerWebhook(sessionWebhook, sessionId, 'message_ack', { message, ack })
|
||||||
|
triggerWebSocket(sessionId, 'message_ack', { message, ack })
|
||||||
if (setMessagesAsSeen) {
|
if (setMessagesAsSeen) {
|
||||||
sendMessageSeenStatus(message)
|
sendMessageSeenStatus(message)
|
||||||
}
|
}
|
||||||
|
|
@ -295,6 +312,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('message_create', async (message) => {
|
client.on('message_create', async (message) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'message_create', { message })
|
triggerWebhook(sessionWebhook, sessionId, 'message_create', { message })
|
||||||
|
triggerWebSocket(sessionId, 'message_create', { message })
|
||||||
if (setMessagesAsSeen) {
|
if (setMessagesAsSeen) {
|
||||||
sendMessageSeenStatus(message)
|
sendMessageSeenStatus(message)
|
||||||
}
|
}
|
||||||
|
|
@ -305,6 +323,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('message_reaction', (reaction) => {
|
client.on('message_reaction', (reaction) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'message_reaction', { reaction })
|
triggerWebhook(sessionWebhook, sessionId, 'message_reaction', { reaction })
|
||||||
|
triggerWebSocket(sessionId, 'message_reaction', { reaction })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -312,6 +331,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('message_edit', (message, newBody, prevBody) => {
|
client.on('message_edit', (message, newBody, prevBody) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'message_edit', { message, newBody, prevBody })
|
triggerWebhook(sessionWebhook, sessionId, 'message_edit', { message, newBody, prevBody })
|
||||||
|
triggerWebSocket(sessionId, 'message_edit', { message, newBody, prevBody })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -319,6 +339,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('message_ciphertext', (message) => {
|
client.on('message_ciphertext', (message) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'message_ciphertext', { message })
|
triggerWebhook(sessionWebhook, sessionId, 'message_ciphertext', { message })
|
||||||
|
triggerWebSocket(sessionId, 'message_ciphertext', { message })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -326,6 +347,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('message_revoke_everyone', async (message) => {
|
client.on('message_revoke_everyone', async (message) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'message_revoke_everyone', { message })
|
triggerWebhook(sessionWebhook, sessionId, 'message_revoke_everyone', { message })
|
||||||
|
triggerWebSocket(sessionId, 'message_revoke_everyone', { message })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -333,6 +355,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('message_revoke_me', async (message, revokedMsg) => {
|
client.on('message_revoke_me', async (message, revokedMsg) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'message_revoke_me', { message, revokedMsg })
|
triggerWebhook(sessionWebhook, sessionId, 'message_revoke_me', { message, revokedMsg })
|
||||||
|
triggerWebSocket(sessionId, 'message_revoke_me', { message, revokedMsg })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -342,6 +365,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
checkIfEventisEnabled('qr')
|
checkIfEventisEnabled('qr')
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'qr', { qr })
|
triggerWebhook(sessionWebhook, sessionId, 'qr', { qr })
|
||||||
|
triggerWebSocket(sessionId, 'qr', { qr })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -349,6 +373,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('ready', () => {
|
client.on('ready', () => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'ready')
|
triggerWebhook(sessionWebhook, sessionId, 'ready')
|
||||||
|
triggerWebSocket(sessionId, 'ready')
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -356,6 +381,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('contact_changed', async (message, oldId, newId, isContact) => {
|
client.on('contact_changed', async (message, oldId, newId, isContact) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'contact_changed', { message, oldId, newId, isContact })
|
triggerWebhook(sessionWebhook, sessionId, 'contact_changed', { message, oldId, newId, isContact })
|
||||||
|
triggerWebSocket(sessionId, 'contact_changed', { message, oldId, newId, isContact })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -363,6 +389,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('chat_removed', async (chat) => {
|
client.on('chat_removed', async (chat) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'chat_removed', { chat })
|
triggerWebhook(sessionWebhook, sessionId, 'chat_removed', { chat })
|
||||||
|
triggerWebSocket(sessionId, 'chat_removed', { chat })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -370,6 +397,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('chat_archived', async (chat, currState, prevState) => {
|
client.on('chat_archived', async (chat, currState, prevState) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'chat_archived', { chat, currState, prevState })
|
triggerWebhook(sessionWebhook, sessionId, 'chat_archived', { chat, currState, prevState })
|
||||||
|
triggerWebSocket(sessionId, 'chat_archived', { chat, currState, prevState })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -377,6 +405,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('unread_count', async (chat) => {
|
client.on('unread_count', async (chat) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'unread_count', { chat })
|
triggerWebhook(sessionWebhook, sessionId, 'unread_count', { chat })
|
||||||
|
triggerWebSocket(sessionId, 'unread_count', { chat })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -384,6 +413,7 @@ const initializeEvents = (client, sessionId) => {
|
||||||
.then(_ => {
|
.then(_ => {
|
||||||
client.on('vote_update', async (vote) => {
|
client.on('vote_update', async (vote) => {
|
||||||
triggerWebhook(sessionWebhook, sessionId, 'vote_update', { vote })
|
triggerWebhook(sessionWebhook, sessionId, 'vote_update', { vote })
|
||||||
|
triggerWebSocket(sessionId, 'vote_update', { vote })
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
@ -447,6 +477,11 @@ const deleteSession = async (sessionId, validation) => {
|
||||||
}
|
}
|
||||||
client.pupPage?.removeAllListeners('close')
|
client.pupPage?.removeAllListeners('close')
|
||||||
client.pupPage?.removeAllListeners('error')
|
client.pupPage?.removeAllListeners('error')
|
||||||
|
try {
|
||||||
|
await terminateWebSocketServer(sessionId)
|
||||||
|
} catch (error) {
|
||||||
|
logger.error({ sessionId, err: error }, 'Failed to terminate WebSocket server')
|
||||||
|
}
|
||||||
if (validation.success) {
|
if (validation.success) {
|
||||||
// Client Connected, request logout
|
// Client Connected, request logout
|
||||||
logger.info({ sessionId }, 'Logging out session')
|
logger.info({ sessionId }, 'Logging out session')
|
||||||
|
|
@ -462,8 +497,8 @@ const deleteSession = async (sessionId, validation) => {
|
||||||
await new Promise(resolve => setTimeout(resolve, 1000))
|
await new Promise(resolve => setTimeout(resolve, 1000))
|
||||||
maxDelay++
|
maxDelay++
|
||||||
}
|
}
|
||||||
await deleteSessionFolder(sessionId)
|
|
||||||
sessions.delete(sessionId)
|
sessions.delete(sessionId)
|
||||||
|
await deleteSessionFolder(sessionId)
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
logger.error({ sessionId, err: error }, 'Failed to delete session')
|
logger.error({ sessionId, err: error }, 'Failed to delete session')
|
||||||
throw error
|
throw error
|
||||||
|
|
|
||||||
|
|
@ -1,12 +1,14 @@
|
||||||
const axios = require('axios')
|
const axios = require('axios')
|
||||||
const { globalApiKey, disabledCallbacks } = require('./config')
|
const { globalApiKey, disabledCallbacks, enableWebHook } = require('./config')
|
||||||
const { logger } = require('./logger')
|
const { logger } = require('./logger')
|
||||||
|
|
||||||
// Trigger webhook endpoint
|
// Trigger webhook endpoint
|
||||||
const triggerWebhook = (webhookURL, sessionId, dataType, data) => {
|
const triggerWebhook = (webhookURL, sessionId, dataType, data) => {
|
||||||
|
if (enableWebHook) {
|
||||||
axios.post(webhookURL, { dataType, data, sessionId }, { headers: { 'x-api-key': globalApiKey } })
|
axios.post(webhookURL, { dataType, data, sessionId }, { headers: { 'x-api-key': globalApiKey } })
|
||||||
.then(() => logger.debug({ sessionId, dataType, data: data || '' }, `New webhook message sent to ${webhookURL}`))
|
.then(() => logger.debug({ sessionId, dataType, data: data || '' }, `New webhook message sent to ${webhookURL}`))
|
||||||
.catch(error => logger.error({ sessionId, dataType, err: error, data: data || '' }, `Failed to send new webhook message to ${webhookURL}`))
|
.catch(error => logger.error({ sessionId, dataType, err: error, data: data || '' }, `Failed to send new webhook message to ${webhookURL}`))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Function to send a response with error status and message
|
// Function to send a response with error status and message
|
||||||
|
|
|
||||||
70
src/websocket.js
Normal file
70
src/websocket.js
Normal file
|
|
@ -0,0 +1,70 @@
|
||||||
|
const { WebSocketServer } = require('ws')
|
||||||
|
const { enableWebSocket } = require('./config')
|
||||||
|
const { logger } = require('./logger')
|
||||||
|
const wssMap = new Map()
|
||||||
|
|
||||||
|
// Function to initialize the WebSocket server if enabled
|
||||||
|
const initWebSocketServer = (sessionId) => {
|
||||||
|
if (enableWebSocket) {
|
||||||
|
const server = wssMap.get(sessionId)
|
||||||
|
if (server) {
|
||||||
|
// happens on session restart
|
||||||
|
return
|
||||||
|
}
|
||||||
|
// init websocket server
|
||||||
|
const wss = new WebSocketServer({ noServer: true })
|
||||||
|
wssMap.set(sessionId, wss)
|
||||||
|
wss.on('connection', (ws) => {
|
||||||
|
logger.debug({ sessionId }, 'WebSocket connection established')
|
||||||
|
ws.on('close', () => {
|
||||||
|
logger.debug({ sessionId }, 'WebSocket connection closed')
|
||||||
|
})
|
||||||
|
ws.on('error', () => {
|
||||||
|
logger.error({ sessionId }, 'WebSocket connection error')
|
||||||
|
})
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Function to initialize the WebSocket server
|
||||||
|
const terminateWebSocketServer = async (sessionId) => {
|
||||||
|
const server = wssMap.get(sessionId)
|
||||||
|
if (!server) {
|
||||||
|
return Promise.resolve()
|
||||||
|
}
|
||||||
|
const closeEventSignal = new Promise((resolve, reject) =>
|
||||||
|
server.close(err => (err ? reject(err) : resolve(undefined)))
|
||||||
|
)
|
||||||
|
for (const ws of server.clients) {
|
||||||
|
ws.terminate()
|
||||||
|
}
|
||||||
|
wssMap.delete(sessionId)
|
||||||
|
await closeEventSignal
|
||||||
|
}
|
||||||
|
|
||||||
|
const triggerWebSocket = (sessionId, dataType, data) => {
|
||||||
|
const server = wssMap.get(sessionId)
|
||||||
|
if (server) {
|
||||||
|
for (const ws of server.clients) {
|
||||||
|
ws.send(JSON.stringify({ dataType, data, sessionId }))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
const handleUpgrade = (request, socket, head) => {
|
||||||
|
const baseUrl = 'ws://' + request.headers.host + '/'
|
||||||
|
const { pathname } = new URL(request.url, baseUrl)
|
||||||
|
if (pathname.startsWith('/ws')) {
|
||||||
|
const sessionId = pathname.split('/')[2]
|
||||||
|
const server = wssMap.get(sessionId)
|
||||||
|
if (server) {
|
||||||
|
server.handleUpgrade(request, socket, head, (ws) => {
|
||||||
|
server.emit('connection', ws, request)
|
||||||
|
})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
socket.destroy()
|
||||||
|
}
|
||||||
|
|
||||||
|
module.exports = { initWebSocketServer, terminateWebSocketServer, handleUpgrade, triggerWebSocket }
|
||||||
Loading…
Reference in a new issue