From 0f0bdcac56e96b8fcab896b65f541b50bf6494ea Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Fri, 21 Aug 2026 11:44:43 +0300 Subject: [PATCH 01/15] increase polling frequency and limit --- src/crons/websocket/websocket.cron.service.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/crons/websocket/websocket.cron.service.ts b/src/crons/websocket/websocket.cron.service.ts index 3e5eafb7c..cd2a207a1 100644 --- a/src/crons/websocket/websocket.cron.service.ts +++ b/src/crons/websocket/websocket.cron.service.ts @@ -145,8 +145,8 @@ export class WebsocketCronService implements OnModuleInit { const stats = await this.networkService.getStats(); - const pollingDelay = stats.refreshRate / 2; - const pollingMaxAttempts = 10; + const pollingDelay = stats.refreshRate / 10; + const pollingMaxAttempts = 15; while (roundToProcessTimestampMs <= latestRoundOnChainData.timestampMs) { await this.pollUntil(async () => await this.isElasticDataAvailableForTimestampMs(roundToProcessTimestampMs, stats), pollingDelay, pollingMaxAttempts); From 3e210d26896d249c697075f7f089f201d77c81f6 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Wed, 26 Aug 2026 14:29:54 +0300 Subject: [PATCH 02/15] fixes + improvements --- .../custom.subscriptions.data.fetcher.ts | 128 ++++++++++++++++++ src/crons/websocket/events.custom.gateway.ts | 29 +--- .../websocket/transaction.custom.gateway.ts | 29 +--- .../websocket/transfers.custom.gateway.ts | 31 +---- src/crons/websocket/websocket.cron.service.ts | 28 ++-- .../websocket.subscription.module.ts | 2 + 6 files changed, 152 insertions(+), 95 deletions(-) create mode 100644 src/crons/websocket/custom.subscriptions.data.fetcher.ts diff --git a/src/crons/websocket/custom.subscriptions.data.fetcher.ts b/src/crons/websocket/custom.subscriptions.data.fetcher.ts new file mode 100644 index 000000000..87cb9002e --- /dev/null +++ b/src/crons/websocket/custom.subscriptions.data.fetcher.ts @@ -0,0 +1,128 @@ +import { Injectable } from '@nestjs/common'; +import { OriginLogger } from '@multiversx/sdk-nestjs-common'; +import { QueryPagination } from 'src/common/entities/query.pagination'; +import { Transaction } from 'src/endpoints/transactions/entities/transaction'; +import { TransactionDetailed } from 'src/endpoints/transactions/entities/transaction.detailed'; +import { TransactionFilter } from 'src/endpoints/transactions/entities/transaction.filter'; +import { TransactionQueryOptions } from 'src/endpoints/transactions/entities/transactions.query.options'; +import { TransactionType } from 'src/endpoints/transactions/entities/transaction.type'; +import { TransferService } from 'src/endpoints/transfers/transfer.service'; +import { EventsService } from 'src/endpoints/events/events.service'; +import { EventsFilter } from 'src/endpoints/events/entities/events.filter'; +import { Events } from 'src/endpoints/events/entities/events'; + +export class CustomSubscriptionsRoundData { + constructor(init?: Partial) { + Object.assign(this, init); + } + + transactions: Transaction[] = []; + transfers: Transaction[] = []; + events: Events[] = []; +} + +@Injectable() +export class CustomSubscriptionsDataFetcher { + private readonly logger = new OriginLogger(CustomSubscriptionsDataFetcher.name); + + private static readonly batchSize = 10000; + + constructor( + private readonly transferService: TransferService, + private readonly eventsService: EventsService, + ) { } + + // Fetches everything the custom subscription gateways need for a given round, once. + // The 'operations' index holds both transactions and smart contract results, so the + // transactions payload is derived from the transfers payload instead of being queried again. + async fetchRoundData(timestampMs: number): Promise { + const [transfers, events] = await Promise.all([ + this.fetchTransfers(timestampMs), + this.fetchEvents(timestampMs), + ]); + + return new CustomSubscriptionsRoundData({ + transfers, + transactions: this.extractTransactions(transfers), + events, + }); + } + + private async fetchTransfers(timestampMs: number): Promise { + try { + const size = CustomSubscriptionsDataFetcher.batchSize; + const filter = new TransactionFilter({ before: timestampMs, after: timestampMs, withTxsRelayedByAddress: true }); + const options = new TransactionQueryOptions({ withScamInfo: false, withUsername: true, withBlockInfo: false, withLogs: false, withOperations: false, withActionTransferValue: false, withTxsOrder: false }); + + const allTransfers: Transaction[] = []; + + let batch = await this.transferService.getTransfers(filter, new QueryPagination({ size }), options); + allTransfers.push(...batch); + + while (batch.length === size) { + const searchAfter = batch[batch.length - 1].searchAfter; + if (searchAfter == null) { + break; + } + + batch = await this.transferService.getTransfers(filter, new QueryPagination({ size, searchAfter }), options); + + allTransfers.push(...batch); + } + + return allTransfers; + } catch (error) { + this.logger.error(`Error fetching transfers for timestamp '${timestampMs}'`); + this.logger.error(error); + return []; + } + } + + private async fetchEvents(timestampMs: number): Promise { + try { + const size = CustomSubscriptionsDataFetcher.batchSize; + const filter = new EventsFilter({ before: timestampMs, after: timestampMs }); + + const allEvents: Events[] = []; + + let batch = await this.eventsService.getEvents(new QueryPagination({ size }), filter); + allEvents.push(...batch); + + while (batch.length === size) { + const searchAfter = batch[batch.length - 1].searchAfter; + if (searchAfter == null) { + break; + } + + batch = await this.eventsService.getEvents(new QueryPagination({ size, searchAfter }), filter); + + allEvents.push(...batch); + } + + return allEvents; + } catch (error) { + this.logger.error(`Error fetching events for timestamp '${timestampMs}'`); + this.logger.error(error); + return []; + } + } + + // Transactions are the 'normal' subset of the operations returned for transfers. They are + // cloned so that clearing the type does not alter the objects broadcast on the transfers channel. + private extractTransactions(transfers: Transaction[]): Transaction[] { + const transactions: Transaction[] = []; + + for (const transfer of transfers) { + if (transfer.type !== TransactionType.Transaction) { + continue; + } + + const transaction = Object.assign(new TransactionDetailed(), transfer); + transaction.type = undefined; + + transactions.push(transaction); + } + + return transactions; + } +} diff --git a/src/crons/websocket/events.custom.gateway.ts b/src/crons/websocket/events.custom.gateway.ts index 7c7406238..9aa35673a 100644 --- a/src/crons/websocket/events.custom.gateway.ts +++ b/src/crons/websocket/events.custom.gateway.ts @@ -2,13 +2,10 @@ import { WebSocketGateway, WebSocketServer, SubscribeMessage, ConnectedSocket, M import { Server, Socket } from 'socket.io'; import { UseFilters, UseInterceptors } from '@nestjs/common'; import { OriginLogger } from '@multiversx/sdk-nestjs-common'; -import { QueryPagination } from 'src/common/entities/query.pagination'; import { WsValidationPipe } from 'src/utils/ws-validation.pipe'; import { WebsocketExceptionsFilter } from 'src/utils/ws-exceptions.filter'; import { RoomKeyGenerator } from './room.key.generator'; -import { EventsService } from 'src/endpoints/events/events.service'; import { EventsCustomSubscribePayload } from 'src/endpoints/events/entities/events.custom.subscribe'; -import { EventsFilter } from 'src/endpoints/events/entities/events.filter'; import { Events } from 'src/endpoints/events/entities/events'; import { LockingGuardInterceptor } from 'src/utils/locking.guard.interceptor'; @@ -22,10 +19,6 @@ export class EventsCustomGateway { @WebSocketServer() server!: Server; - constructor( - private readonly eventsService: EventsService, - ) { } - @UseInterceptors(LockingGuardInterceptor) @SubscribeMessage('subscribeCustomEvents') async handleCustomSubscription( @@ -57,29 +50,11 @@ export class EventsCustomGateway { return { status: 'unsubscribed' }; } - async pushEventsForTimestampMs(timestampMs: number): Promise { + pushEventsForTimestampMs(timestampMs: number, events: Events[]): void { try { - const allEvents: Events[] = []; - const size = 10000; - const filter = new EventsFilter({ before: timestampMs, after: timestampMs }); - - let batch = await this.eventsService.getEvents(new QueryPagination({ size }), filter); - allEvents.push(...batch); - - while (batch.length === size) { - const searchAfter = batch[batch.length - 1].searchAfter; - if (searchAfter == null) { - break; - } - - batch = await this.eventsService.getEvents(new QueryPagination({ size, searchAfter }), filter); - - allEvents.push(...batch); - } - const eventsFilteredForBroadcast: Map = new Map(); - for (const event of allEvents) { + for (const event of events) { const roomKeys = RoomKeyGenerator.generate( EventsCustomGateway.keyPrefix, event, diff --git a/src/crons/websocket/transaction.custom.gateway.ts b/src/crons/websocket/transaction.custom.gateway.ts index 94270cd27..a843f63b6 100644 --- a/src/crons/websocket/transaction.custom.gateway.ts +++ b/src/crons/websocket/transaction.custom.gateway.ts @@ -1,8 +1,5 @@ import { WebSocketGateway, WebSocketServer, SubscribeMessage, ConnectedSocket, MessageBody } from '@nestjs/websockets'; import { Server, Socket } from 'socket.io'; -import { TransactionService } from '../../endpoints/transactions/transaction.service'; -import { TransactionFilter } from '../../endpoints/transactions/entities/transaction.filter'; -import { QueryPagination } from 'src/common/entities/query.pagination'; import { WsValidationPipe } from 'src/utils/ws-validation.pipe'; import { WebsocketExceptionsFilter } from 'src/utils/ws-exceptions.filter'; import { UseFilters, UseInterceptors } from '@nestjs/common'; @@ -20,10 +17,6 @@ export class TransactionsCustomGateway { @WebSocketServer() server!: Server; - constructor( - private readonly transactionService: TransactionService, - ) { } - @UseInterceptors(LockingGuardInterceptor) @SubscribeMessage('subscribeCustomTransactions') async handleCustomSubscription( @@ -52,28 +45,10 @@ export class TransactionsCustomGateway { return { status: 'unsubscribed' }; } - async pushTransactionsForTimestampMs(timestampMs: number): Promise { + pushTransactionsForTimestampMs(timestampMs: number, transactions: Transaction[]): void { try { - const allTransactions: Transaction[] = []; - const size = 10000; - const filter = new TransactionFilter({ before: timestampMs, after: timestampMs }); - - let batch = await this.transactionService.getTransactions(filter, new QueryPagination({ size })); - allTransactions.push(...batch); - - while (batch.length === size) { - const searchAfter = batch[batch.length - 1].searchAfter; - if (searchAfter == null) { - break; - } - - batch = await this.transactionService.getTransactions(filter, new QueryPagination({ size, searchAfter })); - - allTransactions.push(...batch); - } - const txFilteredForBroadcast: Map = new Map(); - for (const transaction of allTransactions) { + for (const transaction of transactions) { const roomKeys = RoomKeyGenerator.generate( TransactionsCustomGateway.keyPrefix, transaction, diff --git a/src/crons/websocket/transfers.custom.gateway.ts b/src/crons/websocket/transfers.custom.gateway.ts index 27fac7c0b..6826c70b8 100644 --- a/src/crons/websocket/transfers.custom.gateway.ts +++ b/src/crons/websocket/transfers.custom.gateway.ts @@ -1,7 +1,5 @@ import { WebSocketGateway, WebSocketServer, SubscribeMessage, ConnectedSocket, MessageBody } from '@nestjs/websockets'; import { Server, Socket } from 'socket.io'; -import { TransactionFilter } from '../../endpoints/transactions/entities/transaction.filter'; -import { QueryPagination } from 'src/common/entities/query.pagination'; import { WsValidationPipe } from 'src/utils/ws-validation.pipe'; import { WebsocketExceptionsFilter } from 'src/utils/ws-exceptions.filter'; import { UseFilters, UseInterceptors } from '@nestjs/common'; @@ -9,9 +7,7 @@ import { OriginLogger } from '@multiversx/sdk-nestjs-common'; import { RoomKeyGenerator } from './room.key.generator'; import { Transaction } from 'src/endpoints/transactions/entities/transaction'; import { LockingGuardInterceptor } from 'src/utils/locking.guard.interceptor'; -import { TransferService } from 'src/endpoints/transfers/transfer.service'; import { TransferCustomSubscribePayload } from 'src/endpoints/websocket/entities/transfers.custom.payload'; -import { TransactionQueryOptions } from 'src/endpoints/transactions/entities/transactions.query.options'; @UseFilters(WebsocketExceptionsFilter) @WebSocketGateway({ cors: { origin: '*' }, path: '/ws/subscription' }) @@ -21,10 +17,6 @@ export class TransfersCustomGateway { @WebSocketServer() server!: Server; - constructor( - private readonly transferService: TransferService, - ) { } - @UseInterceptors(LockingGuardInterceptor) @SubscribeMessage('subscribeCustomTransfers') async handleCustomSubscription( @@ -53,30 +45,11 @@ export class TransfersCustomGateway { return { status: 'unsubscribed' }; } - async pushTransfersForTimestampMs(timestampMs: number): Promise { + pushTransfersForTimestampMs(timestampMs: number, transfers: Transaction[]): void { try { - const allTransfers: Transaction[] = []; - const size = 10000; - const filter = new TransactionFilter({ before: timestampMs, after: timestampMs, withTxsRelayedByAddress: true }); - const options = new TransactionQueryOptions({ withScamInfo: false, withUsername: true, withBlockInfo: false, withLogs: false, withOperations: false, withActionTransferValue: false, withTxsOrder: false }); - - let batch = await this.transferService.getTransfers(filter, new QueryPagination({ size }), options); - allTransfers.push(...batch); - - while (batch.length === size) { - const searchAfter = batch[batch.length - 1].searchAfter; - if (searchAfter == null) { - break; - } - - batch = await this.transferService.getTransfers(filter, new QueryPagination({ size, searchAfter }), options); - - allTransfers.push(...batch); - } - const transfersFilteredForBroadcast: Map = new Map(); - for (const transfer of allTransfers) { + for (const transfer of transfers) { const roomKeys = RoomKeyGenerator.generate( TransfersCustomGateway.keyPrefix, transfer, diff --git a/src/crons/websocket/websocket.cron.service.ts b/src/crons/websocket/websocket.cron.service.ts index cd2a207a1..7897a5825 100644 --- a/src/crons/websocket/websocket.cron.service.ts +++ b/src/crons/websocket/websocket.cron.service.ts @@ -22,6 +22,7 @@ import { ConnectionHandler } from './connection.handler'; import { EventsCustomGateway } from './events.custom.gateway'; import { TransfersCustomGateway } from './transfers.custom.gateway'; import { ApiConfigService } from 'src/common/api-config/api.config.service'; +import { CustomSubscriptionsDataFetcher } from './custom.subscriptions.data.fetcher'; @Injectable() @WebSocketGateway({ cors: { origin: '*' }, path: '/ws/subscription' }) @@ -46,6 +47,7 @@ export class WebsocketCronService implements OnModuleInit { private readonly transfersCustomGateway: TransfersCustomGateway, private readonly apiConfigService: ApiConfigService, private readonly schedulerRegistry: SchedulerRegistry, + private readonly customSubscriptionsDataFetcher: CustomSubscriptionsDataFetcher, ) { } @@ -135,6 +137,7 @@ export class WebsocketCronService implements OnModuleInit { } const latestRoundOnChainData = await this.getLatestRoundOnChainData(); + const statsPromise = this.networkService.getStats(); latestRoundOnChainData.timestampMs = latestRoundOnChainData.timestampMs ?? latestRoundOnChainData.timestamp * 1000; let roundToProcessTimestampMs = await this.cacheService.getOrSetLocal( @@ -143,26 +146,27 @@ export class WebsocketCronService implements OnModuleInit { CacheInfo.WsTimestampMsToProcess().ttl, ); - const stats = await this.networkService.getStats(); + const stats = await statsPromise; const pollingDelay = stats.refreshRate / 10; const pollingMaxAttempts = 15; while (roundToProcessTimestampMs <= latestRoundOnChainData.timestampMs) { await this.pollUntil(async () => await this.isElasticDataAvailableForTimestampMs(roundToProcessTimestampMs, stats), pollingDelay, pollingMaxAttempts); - // call gateways to process logic for custom subscriptions - await Promise.all([ - this.transactionsCustomGateway.pushTransactionsForTimestampMs(roundToProcessTimestampMs), - this.eventsCustomGateway.pushEventsForTimestampMs(roundToProcessTimestampMs), - this.transfersCustomGateway.pushTransfersForTimestampMs(roundToProcessTimestampMs), - ]); + // fetch the round data once, then let each gateway build its own response out of it + const roundData = await this.customSubscriptionsDataFetcher.fetchRoundData(roundToProcessTimestampMs); + + this.transfersCustomGateway.pushTransfersForTimestampMs(roundToProcessTimestampMs, roundData.transfers); + this.eventsCustomGateway.pushEventsForTimestampMs(roundToProcessTimestampMs, roundData.events); + this.transactionsCustomGateway.pushTransactionsForTimestampMs(roundToProcessTimestampMs, roundData.transactions); + roundToProcessTimestampMs += stats.refreshRate; + this.cacheService.setLocal( + CacheInfo.WsTimestampMsToProcess().key, + roundToProcessTimestampMs, + CacheInfo.WsTimestampMsToProcess().ttl, + ); } - this.cacheService.setLocal( - CacheInfo.WsTimestampMsToProcess().key, - roundToProcessTimestampMs, - CacheInfo.WsTimestampMsToProcess().ttl, - ); } @Cron('*/10 * * * * *') diff --git a/src/crons/websocket/websocket.subscription.module.ts b/src/crons/websocket/websocket.subscription.module.ts index dc57125e3..72eced0f6 100644 --- a/src/crons/websocket/websocket.subscription.module.ts +++ b/src/crons/websocket/websocket.subscription.module.ts @@ -21,6 +21,7 @@ import { TransferModule } from 'src/endpoints/transfers/transfer.module'; import { ApiMetricsModule } from 'src/common/metrics/api.metrics.module'; import { EventEmitterModule } from '@nestjs/event-emitter'; import { PersistenceModule } from 'src/common/persistence/persistence.module'; +import { CustomSubscriptionsDataFetcher } from './custom.subscriptions.data.fetcher'; @Module({ imports: [ @@ -39,6 +40,7 @@ import { PersistenceModule } from 'src/common/persistence/persistence.module'; ], providers: [ WebsocketCronService, + CustomSubscriptionsDataFetcher, ConnectionHandler, BlocksGateway, NetworkGateway, From 5142e3add49ed1d07f68f02aa66a35abb08dd085 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Thu, 27 Aug 2026 14:15:28 +0300 Subject: [PATCH 03/15] skip canBeIgnored flagged transfers --- src/crons/websocket/custom.subscriptions.data.fetcher.ts | 4 ++-- src/endpoints/transactions/entities/transaction.ts | 2 ++ .../transactions/entities/transactions.query.options.ts | 3 +++ src/endpoints/transfers/transfer.service.ts | 5 +++++ 4 files changed, 12 insertions(+), 2 deletions(-) diff --git a/src/crons/websocket/custom.subscriptions.data.fetcher.ts b/src/crons/websocket/custom.subscriptions.data.fetcher.ts index 87cb9002e..3562260fa 100644 --- a/src/crons/websocket/custom.subscriptions.data.fetcher.ts +++ b/src/crons/websocket/custom.subscriptions.data.fetcher.ts @@ -52,7 +52,7 @@ export class CustomSubscriptionsDataFetcher { try { const size = CustomSubscriptionsDataFetcher.batchSize; const filter = new TransactionFilter({ before: timestampMs, after: timestampMs, withTxsRelayedByAddress: true }); - const options = new TransactionQueryOptions({ withScamInfo: false, withUsername: true, withBlockInfo: false, withLogs: false, withOperations: false, withActionTransferValue: false, withTxsOrder: false }); + const options = new TransactionQueryOptions({ withScamInfo: false, withUsername: true, withBlockInfo: false, withLogs: false, withOperations: false, withActionTransferValue: false, withTxsOrder: false, withCanBeIgnoredFlag: true }); const allTransfers: Transaction[] = []; @@ -70,7 +70,7 @@ export class CustomSubscriptionsDataFetcher { allTransfers.push(...batch); } - return allTransfers; + return allTransfers.filter((transfer) => transfer.canBeIgnored !== true); } catch (error) { this.logger.error(`Error fetching transfers for timestamp '${timestampMs}'`); this.logger.error(error); diff --git a/src/endpoints/transactions/entities/transaction.ts b/src/endpoints/transactions/entities/transaction.ts index 9eb99e5c2..34808689b 100644 --- a/src/endpoints/transactions/entities/transaction.ts +++ b/src/endpoints/transactions/entities/transaction.ts @@ -122,6 +122,8 @@ export class Transaction { @ApiProperty({ type: String, nullable: true, required: false }) searchAfter?: string | undefined; + canBeIgnored?: boolean | undefined; + getDate(): Date | undefined { if (this.timestamp) { return new Date(this.timestamp * 1000); diff --git a/src/endpoints/transactions/entities/transactions.query.options.ts b/src/endpoints/transactions/entities/transactions.query.options.ts index b12af14cd..caeec80e2 100644 --- a/src/endpoints/transactions/entities/transactions.query.options.ts +++ b/src/endpoints/transactions/entities/transactions.query.options.ts @@ -19,6 +19,9 @@ export class TransactionQueryOptions { withActionTransferValue?: boolean; withTxsOrder?: boolean; + // used only internal + withCanBeIgnoredFlag?: boolean = false; + static applyDefaultOptions(size: number, options: TransactionQueryOptions): TransactionQueryOptions { if (size <= TransactionQueryOptions.SCAM_INFO_MAX_SIZE) { options.withScamInfo = true; diff --git a/src/endpoints/transfers/transfer.service.ts b/src/endpoints/transfers/transfer.service.ts index 41c5b8c81..fc0ba0d14 100644 --- a/src/endpoints/transfers/transfer.service.ts +++ b/src/endpoints/transfers/transfer.service.ts @@ -131,6 +131,11 @@ export class TransferService { for (const elasticOperation of elasticOperations) { const transaction = ApiUtils.mergeObjects(new TransactionDetailed(), elasticOperation); transaction.type = elasticOperation.type === 'normal' ? TransactionType.Transaction : TransactionType.SmartContractResult; + + if (queryOptions.withCanBeIgnoredFlag) { + transaction.canBeIgnored = elasticOperation.canBeIgnored; + } + if (elasticOperation.relayer) { transaction.relayer = elasticOperation.relayer; transaction.isRelayed = true; From d1f5d969f3c2dc5dbe36855138d4ca73d0010872 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Fri, 28 Aug 2026 14:39:19 +0300 Subject: [PATCH 04/15] reduce ws subscriptions broadcast intervbal --- config/config.devnet.yaml | 2 +- config/config.testnet.yaml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/config/config.devnet.yaml b/config/config.devnet.yaml index 96bc5818e..2b829c64a 100644 --- a/config/config.devnet.yaml +++ b/config/config.devnet.yaml @@ -25,7 +25,7 @@ features: port: 6002 maxSubscriptionsPerInstance: 10000 maxSubscriptionsPerClient: 10 - broadcastIntervalMs: 1000 + broadcastIntervalMs: 600 eventsNotifier: enabled: false port: 5674 diff --git a/config/config.testnet.yaml b/config/config.testnet.yaml index f5ab3e1f2..21f3f8c30 100644 --- a/config/config.testnet.yaml +++ b/config/config.testnet.yaml @@ -25,7 +25,7 @@ features: port: 6002 maxSubscriptionsPerInstance: 10000 maxSubscriptionsPerClient: 10 - broadcastIntervalMs: 1000 + broadcastIntervalMs: 600 eventsNotifier: enabled: false port: 5674 From 268413fca88faeb45902bdb2d7547b04eb462b42 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Fri, 28 Aug 2026 14:49:16 +0300 Subject: [PATCH 05/15] performance improvements --- src/crons/websocket/network.gateway.ts | 2 +- src/crons/websocket/websocket.cron.service.ts | 28 +++++++++++-------- .../websocket.subscription.module.ts | 2 -- 3 files changed, 18 insertions(+), 14 deletions(-) diff --git a/src/crons/websocket/network.gateway.ts b/src/crons/websocket/network.gateway.ts index 183520e7d..bf46f1775 100644 --- a/src/crons/websocket/network.gateway.ts +++ b/src/crons/websocket/network.gateway.ts @@ -38,7 +38,7 @@ export class NetworkGateway { async pushStats() { if (this.server.sockets.adapter.rooms.has('statsRoom')) { try { - const stats = await this.networkService.getStats(); + const stats = await this.networkService.getStats(true); this.server.to('statsRoom').emit('statsUpdate', stats); } catch (error) { this.logger.error(error); diff --git a/src/crons/websocket/websocket.cron.service.ts b/src/crons/websocket/websocket.cron.service.ts index 7897a5825..8cc4f046e 100644 --- a/src/crons/websocket/websocket.cron.service.ts +++ b/src/crons/websocket/websocket.cron.service.ts @@ -12,9 +12,7 @@ import { MetricsEvents } from 'src/utils/metrics-events.constants'; import { Server } from 'socket.io'; import { CacheService } from '@multiversx/sdk-nestjs-cache'; import { CacheInfo } from 'src/utils/cache.info'; -import { RoundService } from 'src/endpoints/rounds/round.service'; -import { RoundFilter } from 'src/endpoints/rounds/entities/round.filter'; -import { ElasticQuery, ElasticService, QueryType } from '@multiversx/sdk-nestjs-elastic'; +import { ElasticQuery, ElasticService, ElasticSortOrder, QueryType } from '@multiversx/sdk-nestjs-elastic'; import { NetworkService } from 'src/endpoints/network/network.service'; import { Stats } from 'src/endpoints/network/entities/stats'; import { TransactionsCustomGateway } from './transaction.custom.gateway'; @@ -38,7 +36,6 @@ export class WebsocketCronService implements OnModuleInit { private readonly eventsGateway: EventsGateway, private readonly eventEmitter: EventEmitter2, private readonly cacheService: CacheService, - private readonly roundService: RoundService, private readonly elasticService: ElasticService, private readonly networkService: NetworkService, private readonly transactionsCustomGateway: TransactionsCustomGateway, @@ -136,13 +133,13 @@ export class WebsocketCronService implements OnModuleInit { return; } - const latestRoundOnChainData = await this.getLatestRoundOnChainData(); - const statsPromise = this.networkService.getStats(); - latestRoundOnChainData.timestampMs = latestRoundOnChainData.timestampMs ?? latestRoundOnChainData.timestamp * 1000; + const latestRoundOnChainTimestamp = await this.getLatestRoundOnChainTimestamp(); + const statsPromise = this.networkService.getStats(false); + latestRoundOnChainTimestamp.timestampMs = latestRoundOnChainTimestamp.timestampMs ?? latestRoundOnChainTimestamp.timestamp * 1000; let roundToProcessTimestampMs = await this.cacheService.getOrSetLocal( CacheInfo.WsTimestampMsToProcess().key, - async () => await Promise.resolve(latestRoundOnChainData.timestampMs ?? latestRoundOnChainData.timestamp * 1000), + async () => await Promise.resolve(latestRoundOnChainTimestamp.timestampMs ?? latestRoundOnChainTimestamp.timestamp * 1000), CacheInfo.WsTimestampMsToProcess().ttl, ); @@ -150,7 +147,7 @@ export class WebsocketCronService implements OnModuleInit { const pollingDelay = stats.refreshRate / 10; const pollingMaxAttempts = 15; - while (roundToProcessTimestampMs <= latestRoundOnChainData.timestampMs) { + while (roundToProcessTimestampMs <= latestRoundOnChainTimestamp.timestampMs) { await this.pollUntil(async () => await this.isElasticDataAvailableForTimestampMs(roundToProcessTimestampMs, stats), pollingDelay, pollingMaxAttempts); // fetch the round data once, then let each gateway build its own response out of it @@ -223,8 +220,17 @@ export class WebsocketCronService implements OnModuleInit { }); } - private async getLatestRoundOnChainData() { - const rounds = await this.roundService.getRounds(new RoundFilter({ size: 1 })); + private async getLatestRoundOnChainTimestamp(): Promise<{ timestampMs?: number; timestamp: number }> { + const elasticQuery = ElasticQuery.create(). + withSort([ + { name: 'timestampMs', order: ElasticSortOrder.descending, missing: 0 }, + { name: "timestamp", order: ElasticSortOrder.descending }, + ]) + .withPagination({ from: 0, size: 1 }) + .withFields(['timestampMs', 'timestamp']); + + const rounds = await this.elasticService.getList('rounds', 'round', elasticQuery); + console.log(rounds) return rounds[0]; } diff --git a/src/crons/websocket/websocket.subscription.module.ts b/src/crons/websocket/websocket.subscription.module.ts index 72eced0f6..06ea0ed00 100644 --- a/src/crons/websocket/websocket.subscription.module.ts +++ b/src/crons/websocket/websocket.subscription.module.ts @@ -12,7 +12,6 @@ import { TransactionsGateway } from './transaction.gateway'; import { PoolGateway } from './pool.gateway'; import { EventsGateway } from './events.gateway'; import { ConnectionHandler } from './connection.handler'; -import { RoundModule } from 'src/endpoints/rounds/round.module'; import { TransactionsCustomGateway } from './transaction.custom.gateway'; import { EventsCustomGateway } from './events.custom.gateway'; import { ApiConfigModule } from 'src/common/api-config/api.config.module'; @@ -33,7 +32,6 @@ import { CustomSubscriptionsDataFetcher } from './custom.subscriptions.data.fetc NetworkModule, PoolModule, EventsModule, - RoundModule, TransferModule, ApiConfigModule, ApiMetricsModule, From f460ef343bc49eab6428049d82e43957b50541ea Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Fri, 28 Aug 2026 15:47:02 +0300 Subject: [PATCH 06/15] delete log --- src/crons/websocket/websocket.cron.service.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/crons/websocket/websocket.cron.service.ts b/src/crons/websocket/websocket.cron.service.ts index 8cc4f046e..73cc02c10 100644 --- a/src/crons/websocket/websocket.cron.service.ts +++ b/src/crons/websocket/websocket.cron.service.ts @@ -230,7 +230,7 @@ export class WebsocketCronService implements OnModuleInit { .withFields(['timestampMs', 'timestamp']); const rounds = await this.elasticService.getList('rounds', 'round', elasticQuery); - console.log(rounds) + return rounds[0]; } From da124c273e98e59771e25156dc3211acdf3e7111 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Fri, 28 Aug 2026 15:58:32 +0300 Subject: [PATCH 07/15] add WS messages compression --- config/config.devnet.yaml | 1 + config/config.e2e-mocked.mainnet.yaml | 1 + config/config.e2e.mainnet.yaml | 1 + config/config.mainnet.yaml | 1 + config/config.testnet.yaml | 1 + src/common/api-config/api.config.service.ts | 4 +++ .../websockets/subscription-socket-adapter.ts | 27 +++++++++++++++++++ src/main.ts | 4 +-- 8 files changed, 38 insertions(+), 2 deletions(-) create mode 100644 src/common/websockets/subscription-socket-adapter.ts diff --git a/config/config.devnet.yaml b/config/config.devnet.yaml index 2b829c64a..5dc5508af 100644 --- a/config/config.devnet.yaml +++ b/config/config.devnet.yaml @@ -26,6 +26,7 @@ features: maxSubscriptionsPerInstance: 10000 maxSubscriptionsPerClient: 10 broadcastIntervalMs: 600 + compressionThreshold: 1024 eventsNotifier: enabled: false port: 5674 diff --git a/config/config.e2e-mocked.mainnet.yaml b/config/config.e2e-mocked.mainnet.yaml index 4cf36ee31..02bd8cae4 100644 --- a/config/config.e2e-mocked.mainnet.yaml +++ b/config/config.e2e-mocked.mainnet.yaml @@ -11,6 +11,7 @@ features: maxSubscriptionsPerInstance: 10000 maxSubscriptionsPerClient: 10 broadcastIntervalMs: 6000 + compressionThreshold: 1024 dataApi: enabled: false serviceUrl: 'https://data-api.multiversx.com' diff --git a/config/config.e2e.mainnet.yaml b/config/config.e2e.mainnet.yaml index 3862cd315..5b85e1121 100644 --- a/config/config.e2e.mainnet.yaml +++ b/config/config.e2e.mainnet.yaml @@ -26,6 +26,7 @@ features: maxSubscriptionsPerInstance: 10000 maxSubscriptionsPerClient: 10 broadcastIntervalMs: 6000 + compressionThreshold: 1024 eventsNotifier: enabled: false port: 5674 diff --git a/config/config.mainnet.yaml b/config/config.mainnet.yaml index 6b35bfc0d..faa55d623 100644 --- a/config/config.mainnet.yaml +++ b/config/config.mainnet.yaml @@ -26,6 +26,7 @@ features: maxSubscriptionsPerInstance: 10000 maxSubscriptionsPerClient: 10 broadcastIntervalMs: 6000 + compressionThreshold: 1024 eventsNotifier: enabled: false port: 5674 diff --git a/config/config.testnet.yaml b/config/config.testnet.yaml index 21f3f8c30..0c634e8f6 100644 --- a/config/config.testnet.yaml +++ b/config/config.testnet.yaml @@ -26,6 +26,7 @@ features: maxSubscriptionsPerInstance: 10000 maxSubscriptionsPerClient: 10 broadcastIntervalMs: 600 + compressionThreshold: 1024 eventsNotifier: enabled: false port: 5674 diff --git a/src/common/api-config/api.config.service.ts b/src/common/api-config/api.config.service.ts index 9a4e376d9..733f8c6a4 100644 --- a/src/common/api-config/api.config.service.ts +++ b/src/common/api-config/api.config.service.ts @@ -1016,6 +1016,10 @@ export class ApiConfigService { return this.configService.get('features.websocketSubscription.broadcastIntervalMs') ?? 1000; } + getWebsocketSubscriptionCompressionThreshold(): number | undefined { + return this.configService.get('features.websocketSubscription.compressionThreshold'); + } + getWebsocketMaxSubscriptionsPerInstance(): number { return this.configService.get('features.websocketSubscription.maxSubscriptionsPerInstance') ?? 10_000; } diff --git a/src/common/websockets/subscription-socket-adapter.ts b/src/common/websockets/subscription-socket-adapter.ts new file mode 100644 index 000000000..a89d1a318 --- /dev/null +++ b/src/common/websockets/subscription-socket-adapter.ts @@ -0,0 +1,27 @@ +import { INestApplicationContext } from '@nestjs/common'; +import { IoAdapter } from '@nestjs/platform-socket.io'; +import { ServerOptions } from 'socket.io'; + +// engine.io enables httpCompression by default, but that only covers the long polling transport. +// perMessageDeflate is off unless it is passed explicitly, so websocket frames - which is what +// subscribers actually end up on - go out uncompressed. +export class SubscriptionSocketAdapter extends IoAdapter { + constructor( + app: INestApplicationContext, + private readonly compressionThreshold: number | undefined, + ) { + super(app); + } + + createIOServer(port: number, options?: ServerOptions & { namespace?: string, server?: any }) { + if (this.compressionThreshold === undefined) { + return super.createIOServer(port, options); + } + + // payloads below the threshold are sent as they are, compressing them costs more than it saves + return super.createIOServer(port, { + ...options, + perMessageDeflate: { threshold: this.compressionThreshold }, + }); + } +} diff --git a/src/main.ts b/src/main.ts index fd7f09bfc..c2e65471a 100644 --- a/src/main.ts +++ b/src/main.ts @@ -18,6 +18,7 @@ import { ElasticUpdaterModule } from './crons/elastic.updater/elastic.updater.mo import { PluginService } from './common/plugins/plugin.service'; import { TransactionCompletedModule } from './crons/transaction.processor/transaction.completed.module'; import { SocketAdapter } from './common/websockets/socket-adapter'; +import { SubscriptionSocketAdapter } from './common/websockets/subscription-socket-adapter'; import { ApiConfigModule } from './common/api-config/api.config.module'; import { CacheService, CachingInterceptor, GuestCacheInterceptor, GuestCacheService } from '@multiversx/sdk-nestjs-cache'; import { LoggerInitializer } from '@multiversx/sdk-nestjs-common'; @@ -36,7 +37,6 @@ import { NotWritableError } from './common/indexer/entities/not.writable.error'; import * as bodyParser from 'body-parser'; import * as requestIp from 'request-ip'; import compression from 'compression'; -import { IoAdapter } from '@nestjs/platform-socket.io'; import { WebsocketSubscriptionModule } from './crons/websocket/websocket.subscription.module'; import { RestrictedRoutesMiddleware } from './utils/restricted.routes.middleware'; @@ -92,7 +92,7 @@ async function bootstrap() { if (apiConfigService.getIsWebsocketSubscriptionActive()) { const websocketSubscriptionApp = await NestFactory.create(WebsocketSubscriptionModule); - websocketSubscriptionApp.useWebSocketAdapter(new IoAdapter(websocketSubscriptionApp)); + websocketSubscriptionApp.useWebSocketAdapter(new SubscriptionSocketAdapter(websocketSubscriptionApp, apiConfigService.getWebsocketSubscriptionCompressionThreshold())); await websocketSubscriptionApp.listen(apiConfigService.getWebsocketSubscriptionPort()); } From f0d38916674c79c6d1a503d7e7aae284ca315df7 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Mon, 31 Aug 2026 11:09:21 +0300 Subject: [PATCH 08/15] remove comm --- src/common/websockets/subscription-socket-adapter.ts | 3 --- 1 file changed, 3 deletions(-) diff --git a/src/common/websockets/subscription-socket-adapter.ts b/src/common/websockets/subscription-socket-adapter.ts index a89d1a318..7b51d20f4 100644 --- a/src/common/websockets/subscription-socket-adapter.ts +++ b/src/common/websockets/subscription-socket-adapter.ts @@ -2,9 +2,6 @@ import { INestApplicationContext } from '@nestjs/common'; import { IoAdapter } from '@nestjs/platform-socket.io'; import { ServerOptions } from 'socket.io'; -// engine.io enables httpCompression by default, but that only covers the long polling transport. -// perMessageDeflate is off unless it is passed explicitly, so websocket frames - which is what -// subscribers actually end up on - go out uncompressed. export class SubscriptionSocketAdapter extends IoAdapter { constructor( app: INestApplicationContext, From 44b85b384736bb3ea3f8f128430e7b4d96c9cda6 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Tue, 1 Sep 2026 11:20:29 +0300 Subject: [PATCH 09/15] don't use cache on stats fetch --- src/crons/websocket/websocket.cron.service.ts | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/src/crons/websocket/websocket.cron.service.ts b/src/crons/websocket/websocket.cron.service.ts index 73cc02c10..97b686593 100644 --- a/src/crons/websocket/websocket.cron.service.ts +++ b/src/crons/websocket/websocket.cron.service.ts @@ -133,8 +133,9 @@ export class WebsocketCronService implements OnModuleInit { return; } + const statsPromise = this.networkService.getStats(true); const latestRoundOnChainTimestamp = await this.getLatestRoundOnChainTimestamp(); - const statsPromise = this.networkService.getStats(false); + latestRoundOnChainTimestamp.timestampMs = latestRoundOnChainTimestamp.timestampMs ?? latestRoundOnChainTimestamp.timestamp * 1000; let roundToProcessTimestampMs = await this.cacheService.getOrSetLocal( From a4e875f51065650949cf5d81d7adf12596432339 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Tue, 1 Sep 2026 12:07:19 +0300 Subject: [PATCH 10/15] use deterministic payload stringify on subscribe --- src/crons/websocket/blocks.gateway.ts | 5 +++-- src/crons/websocket/events.gateway.ts | 2 +- src/crons/websocket/transaction.gateway.ts | 2 +- 3 files changed, 5 insertions(+), 4 deletions(-) diff --git a/src/crons/websocket/blocks.gateway.ts b/src/crons/websocket/blocks.gateway.ts index a7df06b7d..fc8be6000 100644 --- a/src/crons/websocket/blocks.gateway.ts +++ b/src/crons/websocket/blocks.gateway.ts @@ -8,6 +8,7 @@ import { UseFilters, UseInterceptors } from '@nestjs/common'; import { WebsocketExceptionsFilter } from 'src/utils/ws-exceptions.filter'; import { WsValidationPipe } from 'src/utils/ws-validation.pipe'; import { OriginLogger } from '@multiversx/sdk-nestjs-common'; +import { RoomKeyGenerator } from './room.key.generator'; import { LockingGuardInterceptor } from 'src/utils/locking.guard.interceptor'; @UseFilters(WebsocketExceptionsFilter) @@ -27,7 +28,7 @@ export class BlocksGateway { @ConnectedSocket() client: Socket, @MessageBody(new WsValidationPipe()) payload: BlockSubscribePayload ) { - const filterIdentifier = JSON.stringify(payload); + const filterIdentifier = RoomKeyGenerator.deterministicStringify(payload); const roomName = `${BlocksGateway.keyPrefix}${filterIdentifier}`; if (!client.rooms.has(roomName)) { @@ -42,7 +43,7 @@ export class BlocksGateway { @ConnectedSocket() client: Socket, @MessageBody(new WsValidationPipe()) payload: BlockSubscribePayload ) { - const filterIdentifier = JSON.stringify(payload); + const filterIdentifier = RoomKeyGenerator.deterministicStringify(payload); const roomName = `${BlocksGateway.keyPrefix}${filterIdentifier}`; if (client.rooms.has(roomName)) { diff --git a/src/crons/websocket/events.gateway.ts b/src/crons/websocket/events.gateway.ts index bb7e8b2d3..019beb6f7 100644 --- a/src/crons/websocket/events.gateway.ts +++ b/src/crons/websocket/events.gateway.ts @@ -28,7 +28,7 @@ export class EventsGateway { @ConnectedSocket() client: Socket, @MessageBody(new WsValidationPipe()) payload: EventsSubscribePayload, ) { - const filterIdentifier = JSON.stringify(payload); + const filterIdentifier = RoomKeyGenerator.deterministicStringify(payload); const roomName = `${EventsGateway.keyPrefix}${filterIdentifier}`; if (!client.rooms.has(roomName)) { await client.join(roomName); diff --git a/src/crons/websocket/transaction.gateway.ts b/src/crons/websocket/transaction.gateway.ts index 016e6db86..ad9129799 100644 --- a/src/crons/websocket/transaction.gateway.ts +++ b/src/crons/websocket/transaction.gateway.ts @@ -47,7 +47,7 @@ export class TransactionsGateway { TransactionFilter.validate(transactionFilter, payload.size || 25); - const filterIdentifier = JSON.stringify(payload); + const filterIdentifier = RoomKeyGenerator.deterministicStringify(payload); const roomName = `${TransactionsGateway.keyPrefix}${filterIdentifier}`; if (!client.rooms.has(roomName)) { From b43257c6107d2d039d450278597b405ecb6ec98f Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Tue, 1 Sep 2026 13:08:19 +0300 Subject: [PATCH 11/15] use deterministic substituion instead of raw replace --- src/crons/websocket/room.key.generator.ts | 11 +++++++ .../websocket/transfers.custom.gateway.ts | 2 +- .../websocket/room.key.generator.spec.ts | 29 +++++++++++++++++++ 3 files changed, 41 insertions(+), 1 deletion(-) diff --git a/src/crons/websocket/room.key.generator.ts b/src/crons/websocket/room.key.generator.ts index 84c95d77a..1983db8cf 100644 --- a/src/crons/websocket/room.key.generator.ts +++ b/src/crons/websocket/room.key.generator.ts @@ -88,6 +88,17 @@ export class RoomKeyGenerator { return rooms; } + // Renaming a field can move it to a different position in the sorted key, so the key is + // rebuilt from its parsed form instead of being patched in place. + public static substitute(prefix: string, roomKey: string, from: string, to: string): string { + const { [from]: value, ...rest } = JSON.parse(roomKey.slice(prefix.length)); + if (value === undefined) { + return roomKey; + } + + return `${prefix}${this.deterministicStringify({ ...rest, [to]: value })}`; + } + static deterministicStringify(obj: Record): string { return JSON.stringify( Object.keys(obj) diff --git a/src/crons/websocket/transfers.custom.gateway.ts b/src/crons/websocket/transfers.custom.gateway.ts index 6826c70b8..3a5e47071 100644 --- a/src/crons/websocket/transfers.custom.gateway.ts +++ b/src/crons/websocket/transfers.custom.gateway.ts @@ -60,7 +60,7 @@ export class TransfersCustomGateway { const substitutions = TransferCustomSubscribePayload.getFieldsSubstitutions(); for (const [key, substituteFields] of Object.entries(substitutions)) { for (const substituteField of substituteFields) { - const substituteRoomKey = roomKey.replace(`"${substituteField}":`, `"${key}":`); + const substituteRoomKey = RoomKeyGenerator.substitute(TransfersCustomGateway.keyPrefix, roomKey, substituteField, key); if (this.server.sockets.adapter.rooms.has(substituteRoomKey)) { if (!transfersFilteredForBroadcast.has(substituteRoomKey)) { transfersFilteredForBroadcast.set(substituteRoomKey, []); diff --git a/src/test/unit/crons/websocket/room.key.generator.spec.ts b/src/test/unit/crons/websocket/room.key.generator.spec.ts index 85a0ca4a6..db7c92d9d 100644 --- a/src/test/unit/crons/websocket/room.key.generator.spec.ts +++ b/src/test/unit/crons/websocket/room.key.generator.spec.ts @@ -99,4 +99,33 @@ describe('RoomKeyGenerator', () => { expect(rooms[0]).toBe('{"sender":"alice"}'); }); }); + + describe('substitute', () => { + it('renames a field and keeps the key sorted', () => { + const roomKey = 'p-' + RoomKeyGenerator.deterministicStringify({ function: 'swap', sender: 'alice' }); + + expect(RoomKeyGenerator.substitute('p-', roomKey, 'sender', 'address')) + .toBe('p-' + RoomKeyGenerator.deterministicStringify({ address: 'alice', function: 'swap' })); + }); + + // A plain string replace left the renamed field where the old one sorted, so this combination + // never matched the room the subscriber had actually joined. + it('matches the room key a subscriber with that field would have joined', () => { + const subscribed = 'p-' + RoomKeyGenerator.deterministicStringify({ address: 'alice', function: 'swap' }); + + const generated = RoomKeyGenerator.generate( + 'p-', + { sender: 'alice', function: 'swap' }, + TransactionCustomSubscribePayload, + ).map((roomKey) => RoomKeyGenerator.substitute('p-', roomKey, 'sender', 'address')); + + expect(generated).toContain(subscribed); + }); + + it('leaves the key untouched when the field is not part of it', () => { + const roomKey = 'p-' + RoomKeyGenerator.deterministicStringify({ receiver: 'bob' }); + + expect(RoomKeyGenerator.substitute('p-', roomKey, 'sender', 'address')).toBe(roomKey); + }); + }); }); From ce6cb3f460b7c85e3a45100ab54e42ce2ce90ee5 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Mon, 7 Sep 2026 10:12:29 +0300 Subject: [PATCH 12/15] disable tx pool --- src/common/api-config/api.config.service.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/src/common/api-config/api.config.service.ts b/src/common/api-config/api.config.service.ts index 733f8c6a4..e17af0aac 100644 --- a/src/common/api-config/api.config.service.ts +++ b/src/common/api-config/api.config.service.ts @@ -721,6 +721,8 @@ export class ApiConfigService { } isTransactionPoolEnabled(): boolean { + return false; + //@ts-ignore return this.configService.get('features.transactionPool.enabled') ?? false; } From 8b687c26555ed259070e9d886854dd969bb24309 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Mon, 7 Sep 2026 10:33:34 +0300 Subject: [PATCH 13/15] disable pool in warmer --- src/crons/cache.warmer/cache.warmer.service.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/src/crons/cache.warmer/cache.warmer.service.ts b/src/crons/cache.warmer/cache.warmer.service.ts index e50c8b391..18270661b 100644 --- a/src/crons/cache.warmer/cache.warmer.service.ts +++ b/src/crons/cache.warmer/cache.warmer.service.ts @@ -151,6 +151,8 @@ export class CacheWarmerService { @Lock({ name: 'Transaction pool invalidation', verbose: true }) async handleTxPoolInvalidations() { + return; + //@ts-ignore const pool = await this.poolService.getTxPoolRaw(); await this.invalidateKey(CacheInfo.TransactionPool.key, pool, this.apiConfigService.getTransactionPoolCacheWarmerTtlInSeconds()); From 2fdf4a71d7d6e72e0701743aa811ade547dd6b22 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Mon, 7 Sep 2026 12:42:55 +0300 Subject: [PATCH 14/15] enable tx pool --- src/common/api-config/api.config.service.ts | 2 -- src/crons/cache.warmer/cache.warmer.service.ts | 2 -- 2 files changed, 4 deletions(-) diff --git a/src/common/api-config/api.config.service.ts b/src/common/api-config/api.config.service.ts index e17af0aac..733f8c6a4 100644 --- a/src/common/api-config/api.config.service.ts +++ b/src/common/api-config/api.config.service.ts @@ -721,8 +721,6 @@ export class ApiConfigService { } isTransactionPoolEnabled(): boolean { - return false; - //@ts-ignore return this.configService.get('features.transactionPool.enabled') ?? false; } diff --git a/src/crons/cache.warmer/cache.warmer.service.ts b/src/crons/cache.warmer/cache.warmer.service.ts index 18270661b..e50c8b391 100644 --- a/src/crons/cache.warmer/cache.warmer.service.ts +++ b/src/crons/cache.warmer/cache.warmer.service.ts @@ -151,8 +151,6 @@ export class CacheWarmerService { @Lock({ name: 'Transaction pool invalidation', verbose: true }) async handleTxPoolInvalidations() { - return; - //@ts-ignore const pool = await this.poolService.getTxPoolRaw(); await this.invalidateKey(CacheInfo.TransactionPool.key, pool, this.apiConfigService.getTransactionPoolCacheWarmerTtlInSeconds()); From 05cbffcea94342fecf59738bc262bf46ca99c1f3 Mon Sep 17 00:00:00 2001 From: GuticaStefan Date: Thu, 24 Sep 2026 11:59:05 +0300 Subject: [PATCH 15/15] exclude sockets own rooms from global subscriptions limit --- .../utils/locking.guard.interceptor.spec.ts | 80 +++++++++++++++++++ src/utils/locking.guard.interceptor.ts | 6 +- 2 files changed, 84 insertions(+), 2 deletions(-) create mode 100644 src/test/unit/utils/locking.guard.interceptor.spec.ts diff --git a/src/test/unit/utils/locking.guard.interceptor.spec.ts b/src/test/unit/utils/locking.guard.interceptor.spec.ts new file mode 100644 index 000000000..7724af66a --- /dev/null +++ b/src/test/unit/utils/locking.guard.interceptor.spec.ts @@ -0,0 +1,80 @@ +import { CallHandler, ExecutionContext } from '@nestjs/common'; +import { lastValueFrom, of } from 'rxjs'; +import { ApiConfigService } from 'src/common/api-config/api.config.service'; +import { LockingGuardInterceptor } from 'src/utils/locking.guard.interceptor'; + +describe('LockingGuardInterceptor', () => { + const maxGlobal = 3; + const maxClient = 2; + + const apiConfigService = { + getWebsocketMaxSubscriptionsPerInstance: () => maxGlobal, + getWebsocketMaxSubscriptionsPerClient: () => maxClient, + } as unknown as ApiConfigService; + + const next: CallHandler = { handle: () => of({ status: 'success' }) }; + + // mimics socket.io: every connected socket sits in a room named after its own id + const createServer = (socketIds: string[], subscriptions: Record) => { + const rooms = new Map>(); + for (const id of socketIds) { + rooms.set(id, new Set([id])); + } + + for (const [room, members] of Object.entries(subscriptions)) { + rooms.set(room, new Set(members)); + } + + const sockets = new Map(socketIds.map(id => [id, {}])); + + return { sockets: { adapter: { rooms }, sockets } }; + }; + + const createContext = (server: any, clientId: string): ExecutionContext => { + const clientRooms = new Set(); + for (const [room, members] of server.sockets.adapter.rooms) { + if (members.has(clientId)) { + clientRooms.add(room); + } + } + + const client = { id: clientId, rooms: clientRooms, nsp: { server } }; + + return { switchToWs: () => ({ getClient: () => client }) } as unknown as ExecutionContext; + }; + + it('does not count connected sockets as subscriptions', async () => { + const socketIds = Array.from({ length: 10 }, (_, i) => `socket-${i}`); + const server = createServer(socketIds, { 'tx-{}': ['socket-1'] }); + + const interceptor = new LockingGuardInterceptor(apiConfigService); + const result = await lastValueFrom(interceptor.intercept(createContext(server, 'socket-0'), next)); + + expect(result).toEqual({ status: 'success' }); + }); + + it('rejects when the instance subscriptions limit is reached', async () => { + const server = createServer(['socket-0', 'socket-1'], { + 'tx-{"size":1}': ['socket-1'], + 'tx-{"size":2}': ['socket-1'], + 'blocks-{}': ['socket-1'], + }); + + const interceptor = new LockingGuardInterceptor(apiConfigService); + + await expect(lastValueFrom(interceptor.intercept(createContext(server, 'socket-0'), next))) + .rejects.toThrow(`Max global subscriptions (${maxGlobal}) reached!`); + }); + + it('rejects when the client subscriptions limit is reached', async () => { + const server = createServer(['socket-0'], { + 'tx-{}': ['socket-0'], + 'blocks-{}': ['socket-0'], + }); + + const interceptor = new LockingGuardInterceptor(apiConfigService); + + await expect(lastValueFrom(interceptor.intercept(createContext(server, 'socket-0'), next))) + .rejects.toThrow(`Max client subscriptions (${maxClient}) reached!`); + }); +}); diff --git a/src/utils/locking.guard.interceptor.ts b/src/utils/locking.guard.interceptor.ts index 75e980504..a6d59d050 100644 --- a/src/utils/locking.guard.interceptor.ts +++ b/src/utils/locking.guard.interceptor.ts @@ -34,12 +34,14 @@ export class LockingGuardInterceptor implements NestInterceptor { return from(mutex.acquire()).pipe( switchMap((release) => { try { - const totalRoomsGlobal = client.nsp.server.sockets.adapter.rooms.size; + // every socket is also in a private room named after its id, those are not subscriptions + const namespace = client.nsp.server.sockets; + const totalSubscriptionsGlobal = namespace.adapter.rooms.size - namespace.sockets.size; const totalClientRooms = client.rooms.size; const maxGlobal = this.apiConfigService.getWebsocketMaxSubscriptionsPerInstance(); const maxClient = this.apiConfigService.getWebsocketMaxSubscriptionsPerClient(); - if (totalRoomsGlobal >= maxGlobal) { + if (totalSubscriptionsGlobal >= maxGlobal) { throw new WsException(`Max global subscriptions (${maxGlobal}) reached!`); }