From 3da05b805cb8ac9425aee94d39eb727b2a8dffc1 Mon Sep 17 00:00:00 2001 From: Diego Sampaio Date: Fri, 14 Feb 2025 17:14:08 -0300 Subject: [PATCH] fix: ddp-streamer restart keeping opened connections (#35195) --- _templates/service/new/service.ejs.t | 4 +- apps/meteor/ee/server/startup/index.ts | 4 +- ee/apps/account-service/src/service.ts | 4 +- ee/apps/authorization-service/src/service.ts | 4 +- ee/apps/ddp-streamer/src/service.ts | 11 +- ee/apps/omnichannel-transcript/src/service.ts | 4 +- ee/apps/presence-service/src/service.ts | 4 +- ee/apps/queue-worker/src/service.ts | 4 +- ee/apps/stream-hub-service/src/service.ts | 4 +- ee/packages/network-broker/src/index.ts | 152 +++++++++--------- 10 files changed, 103 insertions(+), 92 deletions(-) diff --git a/_templates/service/new/service.ejs.t b/_templates/service/new/service.ejs.t index bdd160b43cdfc..64b363fe29a29 100644 --- a/_templates/service/new/service.ejs.t +++ b/_templates/service/new/service.ejs.t @@ -2,7 +2,7 @@ to: ee/apps/<%= name %>/src/service.ts --- import { api, getConnection, getTrashCollection } from '@rocket.chat/core-services'; -import { broker } from '@rocket.chat/network-broker'; +import { startBroker } from '@rocket.chat/network-broker'; import { startTracing } from '@rocket.chat/tracing'; import polka from 'polka'; @@ -17,7 +17,7 @@ const PORT = process.env.PORT || <%= h.random() %>; registerServiceModels(db, await getTrashCollection()); - api.setBroker(broker); + api.setBroker(startBroker()); // need to import service after models are registered const { <%= h.changeCase.pascalCase(name) %> } = await import('./<%= h.changeCase.pascalCase(name) %>'); diff --git a/apps/meteor/ee/server/startup/index.ts b/apps/meteor/ee/server/startup/index.ts index 51bb011828ba2..e295467e117cd 100644 --- a/apps/meteor/ee/server/startup/index.ts +++ b/apps/meteor/ee/server/startup/index.ts @@ -12,9 +12,9 @@ import { isRunningMs } from '../../../server/lib/isRunningMs'; export const registerEEBroker = async (): Promise => { // only starts network broker if running in micro services mode if (isRunningMs()) { - const { broker } = await import('@rocket.chat/network-broker'); + const { startBroker } = await import('@rocket.chat/network-broker'); - api.setBroker(broker); + api.setBroker(startBroker()); void api.start(); } else { require('./presence'); diff --git a/ee/apps/account-service/src/service.ts b/ee/apps/account-service/src/service.ts index e9e8cd8bef832..1aa049ba95640 100755 --- a/ee/apps/account-service/src/service.ts +++ b/ee/apps/account-service/src/service.ts @@ -1,5 +1,5 @@ import { api, getConnection, getTrashCollection } from '@rocket.chat/core-services'; -import { broker } from '@rocket.chat/network-broker'; +import { startBroker } from '@rocket.chat/network-broker'; import { startTracing } from '@rocket.chat/tracing'; import polka from 'polka'; @@ -14,7 +14,7 @@ const PORT = process.env.PORT || 3033; registerServiceModels(db, await getTrashCollection()); - api.setBroker(broker); + api.setBroker(startBroker()); // need to import service after models are registered const { Account } = await import('./Account'); diff --git a/ee/apps/authorization-service/src/service.ts b/ee/apps/authorization-service/src/service.ts index 8b3350c5f7544..cd8980459ee69 100755 --- a/ee/apps/authorization-service/src/service.ts +++ b/ee/apps/authorization-service/src/service.ts @@ -1,5 +1,5 @@ import { api, getConnection, getTrashCollection } from '@rocket.chat/core-services'; -import { broker } from '@rocket.chat/network-broker'; +import { startBroker } from '@rocket.chat/network-broker'; import { startTracing } from '@rocket.chat/tracing'; import polka from 'polka'; @@ -14,7 +14,7 @@ const PORT = process.env.PORT || 3034; registerServiceModels(db, await getTrashCollection()); - api.setBroker(broker); + api.setBroker(startBroker()); // need to import service after models are registered const { Authorization } = await import('../../../../apps/meteor/server/services/authorization/service'); diff --git a/ee/apps/ddp-streamer/src/service.ts b/ee/apps/ddp-streamer/src/service.ts index 66ad72be40501..35012c2b8af8b 100755 --- a/ee/apps/ddp-streamer/src/service.ts +++ b/ee/apps/ddp-streamer/src/service.ts @@ -1,5 +1,8 @@ +import os from 'os'; + import { api, getConnection, getTrashCollection } from '@rocket.chat/core-services'; -import { broker } from '@rocket.chat/network-broker'; +import { InstanceStatus } from '@rocket.chat/instance-status'; +import { startBroker } from '@rocket.chat/network-broker'; import { startTracing } from '@rocket.chat/tracing'; import { registerServiceModels } from '../../../../apps/meteor/ee/server/lib/registerServiceModels'; @@ -11,7 +14,11 @@ import { registerServiceModels } from '../../../../apps/meteor/ee/server/lib/reg registerServiceModels(db, await getTrashCollection()); - api.setBroker(broker); + api.setBroker( + startBroker({ + nodeID: `${os.hostname().toLowerCase()}-${InstanceStatus.id()}`, + }), + ); // need to import service after models are registered const { NotificationsModule } = await import('../../../../apps/meteor/server/modules/notifications/notifications.module'); diff --git a/ee/apps/omnichannel-transcript/src/service.ts b/ee/apps/omnichannel-transcript/src/service.ts index e78d646bf9ea9..c50d9493579ca 100644 --- a/ee/apps/omnichannel-transcript/src/service.ts +++ b/ee/apps/omnichannel-transcript/src/service.ts @@ -1,6 +1,6 @@ import { api, getConnection, getTrashCollection } from '@rocket.chat/core-services'; import { Logger } from '@rocket.chat/logger'; -import { broker } from '@rocket.chat/network-broker'; +import { startBroker } from '@rocket.chat/network-broker'; import { startTracing } from '@rocket.chat/tracing'; import polka from 'polka'; @@ -15,7 +15,7 @@ const PORT = process.env.PORT || 3036; registerServiceModels(db, await getTrashCollection()); - api.setBroker(broker); + api.setBroker(startBroker()); // need to import service after models are registered const { OmnichannelTranscript } = await import('@rocket.chat/omnichannel-services'); diff --git a/ee/apps/presence-service/src/service.ts b/ee/apps/presence-service/src/service.ts index f1e86419fc695..f5a2a339afa66 100755 --- a/ee/apps/presence-service/src/service.ts +++ b/ee/apps/presence-service/src/service.ts @@ -1,5 +1,5 @@ import { api, getConnection, getTrashCollection } from '@rocket.chat/core-services'; -import { broker } from '@rocket.chat/network-broker'; +import { startBroker } from '@rocket.chat/network-broker'; import { startTracing } from '@rocket.chat/tracing'; import polka from 'polka'; @@ -14,7 +14,7 @@ const PORT = process.env.PORT || 3031; registerServiceModels(db, await getTrashCollection()); - api.setBroker(broker); + api.setBroker(startBroker()); // need to import Presence service after models are registered const { Presence } = await import('@rocket.chat/presence'); diff --git a/ee/apps/queue-worker/src/service.ts b/ee/apps/queue-worker/src/service.ts index e936621e3389c..ec8a8a06af2e5 100644 --- a/ee/apps/queue-worker/src/service.ts +++ b/ee/apps/queue-worker/src/service.ts @@ -1,6 +1,6 @@ import { api, getConnection, getTrashCollection } from '@rocket.chat/core-services'; import { Logger } from '@rocket.chat/logger'; -import { broker } from '@rocket.chat/network-broker'; +import { startBroker } from '@rocket.chat/network-broker'; import { startTracing } from '@rocket.chat/tracing'; import polka from 'polka'; @@ -15,7 +15,7 @@ const PORT = process.env.PORT || 3038; registerServiceModels(db, await getTrashCollection()); - api.setBroker(broker); + api.setBroker(startBroker()); // need to import service after models are registeredpackagfe const { QueueWorker } = await import('@rocket.chat/omnichannel-services'); diff --git a/ee/apps/stream-hub-service/src/service.ts b/ee/apps/stream-hub-service/src/service.ts index 07baf5923f013..568f143054643 100755 --- a/ee/apps/stream-hub-service/src/service.ts +++ b/ee/apps/stream-hub-service/src/service.ts @@ -1,6 +1,6 @@ import { api, getConnection, getTrashCollection } from '@rocket.chat/core-services'; import { Logger } from '@rocket.chat/logger'; -import { broker } from '@rocket.chat/network-broker'; +import { startBroker } from '@rocket.chat/network-broker'; import { startTracing } from '@rocket.chat/tracing'; import polka from 'polka'; @@ -17,7 +17,7 @@ const PORT = process.env.PORT || 3035; registerServiceModels(db, await getTrashCollection()); - api.setBroker(broker); + api.setBroker(startBroker()); // TODO having to import Logger to pass as a param is a temporary solution. logger should come from the service (either from broker or api) const watcher = new DatabaseWatcher({ db, logger: Logger }); diff --git a/ee/packages/network-broker/src/index.ts b/ee/packages/network-broker/src/index.ts index 4424c084bdaa6..2c6da09750e31 100644 --- a/ee/packages/network-broker/src/index.ts +++ b/ee/packages/network-broker/src/index.ts @@ -1,5 +1,6 @@ import { isMeteorError, MeteorError } from '@rocket.chat/core-services'; import EJSON from 'ejson'; +import type Moleculer from 'moleculer'; import { Errors, Serializers, ServiceBroker } from 'moleculer'; import { pino } from 'pino'; @@ -70,82 +71,85 @@ class EJSONSerializer extends Base { } } -const network = new ServiceBroker({ - namespace: MS_NAMESPACE, - skipProcessEventRegistration: SKIP_PROCESS_EVENT_REGISTRATION === 'true', - transporter: TRANSPORTER, - metrics: { - enabled: MS_METRICS === 'true', - reporter: [ - { - type: 'Prometheus', - options: { - port: MS_METRICS_PORT, +export function startBroker(options: Moleculer.BrokerOptions = {}): NetworkBroker { + const network = new ServiceBroker({ + namespace: MS_NAMESPACE, + skipProcessEventRegistration: SKIP_PROCESS_EVENT_REGISTRATION === 'true', + transporter: TRANSPORTER, + metrics: { + enabled: MS_METRICS === 'true', + reporter: [ + { + type: 'Prometheus', + options: { + port: MS_METRICS_PORT, + }, }, - }, - ], - }, - cacher: CACHE, - serializer: SERIALIZER === 'EJSON' ? new EJSONSerializer() : SERIALIZER, - logger: { - type: 'Pino', - options: { - level: MOLECULER_LOG_LEVEL, - pino: { - options: { - timestamp: pino.stdTimeFunctions.isoTime, - ...(process.env.NODE_ENV !== 'production' - ? { - transport: { - target: 'pino-pretty', - options: { - colorize: true, + ], + }, + cacher: CACHE, + serializer: SERIALIZER === 'EJSON' ? new EJSONSerializer() : SERIALIZER, + logger: { + type: 'Pino', + options: { + level: MOLECULER_LOG_LEVEL, + pino: { + options: { + timestamp: pino.stdTimeFunctions.isoTime, + ...(process.env.NODE_ENV !== 'production' + ? { + transport: { + target: 'pino-pretty', + options: { + colorize: true, + }, }, - }, - } - : {}), + } + : {}), + }, }, }, }, - }, - registry: { - strategy: BALANCE_STRATEGY, - preferLocal: BALANCE_PREFER_LOCAL !== 'false', - }, - - requestTimeout: parseInt(REQUEST_TIMEOUT) * 1000, - retryPolicy: { - enabled: RETRY_ENABLED === 'true', - retries: parseInt(RETRY_RETRIES), - delay: parseInt(RETRY_DELAY), - maxDelay: parseInt(RETRY_MAX_DELAY), - factor: parseInt(RETRY_FACTOR), - check: (err: any): boolean => err && !!err.retryable, - }, - - maxCallLevel: 100, - heartbeatInterval: parseInt(HEARTBEAT_INTERVAL), - heartbeatTimeout: parseInt(HEARTBEAT_TIMEOUT), - - // circuitBreaker: { - // enabled: false, - // threshold: 0.5, - // windowTime: 60, - // minRequestCount: 20, - // halfOpenTime: 10 * 1000, - // check: (err: any): boolean => err && err.code >= 500, - // }, - - bulkhead: { - enabled: BULKHEAD_ENABLED === 'true', - concurrency: parseInt(BULKHEAD_CONCURRENCY), - maxQueueSize: parseInt(BULKHEAD_MAX_QUEUE_SIZE), - }, - - errorRegenerator: new CustomRegenerator(), - started(): void { - console.log('NetworkBroker started successfully.'); - }, -}); - -export const broker = new NetworkBroker(network); + registry: { + strategy: BALANCE_STRATEGY, + preferLocal: BALANCE_PREFER_LOCAL !== 'false', + }, + + requestTimeout: parseInt(REQUEST_TIMEOUT) * 1000, + retryPolicy: { + enabled: RETRY_ENABLED === 'true', + retries: parseInt(RETRY_RETRIES), + delay: parseInt(RETRY_DELAY), + maxDelay: parseInt(RETRY_MAX_DELAY), + factor: parseInt(RETRY_FACTOR), + check: (err: any): boolean => err && !!err.retryable, + }, + + maxCallLevel: 100, + heartbeatInterval: parseInt(HEARTBEAT_INTERVAL), + heartbeatTimeout: parseInt(HEARTBEAT_TIMEOUT), + + // circuitBreaker: { + // enabled: false, + // threshold: 0.5, + // windowTime: 60, + // minRequestCount: 20, + // halfOpenTime: 10 * 1000, + // check: (err: any): boolean => err && err.code >= 500, + // }, + + bulkhead: { + enabled: BULKHEAD_ENABLED === 'true', + concurrency: parseInt(BULKHEAD_CONCURRENCY), + maxQueueSize: parseInt(BULKHEAD_MAX_QUEUE_SIZE), + }, + + errorRegenerator: new CustomRegenerator(), + started(): void { + console.log('NetworkBroker started successfully.'); + }, + ...options, + }); + + return new NetworkBroker(network); +}