diff --git a/config/config.devnet.yaml b/config/config.devnet.yaml index 96bc5818e..5dc5508af 100644 --- a/config/config.devnet.yaml +++ b/config/config.devnet.yaml @@ -25,7 +25,8 @@ features: port: 6002 maxSubscriptionsPerInstance: 10000 maxSubscriptionsPerClient: 10 - broadcastIntervalMs: 1000 + 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 5dbbe43a4..d4d664cbb 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: 600 + 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 df7b38c63..6b3743d12 100644 --- a/config/config.e2e.mainnet.yaml +++ b/config/config.e2e.mainnet.yaml @@ -26,6 +26,7 @@ features: maxSubscriptionsPerInstance: 10000 maxSubscriptionsPerClient: 10 broadcastIntervalMs: 600 + compressionThreshold: 1024 eventsNotifier: enabled: false port: 5674 diff --git a/config/config.mainnet.yaml b/config/config.mainnet.yaml index 9fc343363..618c67e6d 100644 --- a/config/config.mainnet.yaml +++ b/config/config.mainnet.yaml @@ -26,6 +26,7 @@ features: maxSubscriptionsPerInstance: 10000 maxSubscriptionsPerClient: 10 broadcastIntervalMs: 600 + compressionThreshold: 1024 eventsNotifier: enabled: false port: 5674 diff --git a/config/config.testnet.yaml b/config/config.testnet.yaml index f5ab3e1f2..0c634e8f6 100644 --- a/config/config.testnet.yaml +++ b/config/config.testnet.yaml @@ -25,7 +25,8 @@ features: port: 6002 maxSubscriptionsPerInstance: 10000 maxSubscriptionsPerClient: 10 - broadcastIntervalMs: 1000 + 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..7b51d20f4 --- /dev/null +++ b/src/common/websockets/subscription-socket-adapter.ts @@ -0,0 +1,24 @@ +import { INestApplicationContext } from '@nestjs/common'; +import { IoAdapter } from '@nestjs/platform-socket.io'; +import { ServerOptions } from 'socket.io'; + +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/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/custom.subscriptions.data.fetcher.ts b/src/crons/websocket/custom.subscriptions.data.fetcher.ts new file mode 100644 index 000000000..3562260fa --- /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, withCanBeIgnoredFlag: true }); + + 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.filter((transfer) => transfer.canBeIgnored !== true); + } 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/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/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/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/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/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)) { diff --git a/src/crons/websocket/transfers.custom.gateway.ts b/src/crons/websocket/transfers.custom.gateway.ts index 27fac7c0b..3a5e47071 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, @@ -87,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/crons/websocket/websocket.cron.service.ts b/src/crons/websocket/websocket.cron.service.ts index 3e5eafb7c..97b686593 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'; @@ -22,6 +20,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' }) @@ -37,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, @@ -46,6 +44,7 @@ export class WebsocketCronService implements OnModuleInit { private readonly transfersCustomGateway: TransfersCustomGateway, private readonly apiConfigService: ApiConfigService, private readonly schedulerRegistry: SchedulerRegistry, + private readonly customSubscriptionsDataFetcher: CustomSubscriptionsDataFetcher, ) { } @@ -134,35 +133,38 @@ export class WebsocketCronService implements OnModuleInit { return; } - const latestRoundOnChainData = await this.getLatestRoundOnChainData(); - latestRoundOnChainData.timestampMs = latestRoundOnChainData.timestampMs ?? latestRoundOnChainData.timestamp * 1000; + const statsPromise = this.networkService.getStats(true); + const latestRoundOnChainTimestamp = await this.getLatestRoundOnChainTimestamp(); + + 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, ); - const stats = await this.networkService.getStats(); + const stats = await statsPromise; - const pollingDelay = stats.refreshRate / 2; - const pollingMaxAttempts = 10; - while (roundToProcessTimestampMs <= latestRoundOnChainData.timestampMs) { + const pollingDelay = stats.refreshRate / 10; + const pollingMaxAttempts = 15; + while (roundToProcessTimestampMs <= latestRoundOnChainTimestamp.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 * * * * *') @@ -219,8 +221,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); + return rounds[0]; } diff --git a/src/crons/websocket/websocket.subscription.module.ts b/src/crons/websocket/websocket.subscription.module.ts index dc57125e3..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'; @@ -21,6 +20,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: [ @@ -32,13 +32,13 @@ import { PersistenceModule } from 'src/common/persistence/persistence.module'; NetworkModule, PoolModule, EventsModule, - RoundModule, TransferModule, ApiConfigModule, ApiMetricsModule, ], providers: [ WebsocketCronService, + CustomSubscriptionsDataFetcher, ConnectionHandler, BlocksGateway, NetworkGateway, 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 0a51dfb9b..6adf11139 100644 --- a/src/endpoints/transfers/transfer.service.ts +++ b/src/endpoints/transfers/transfer.service.ts @@ -128,6 +128,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; 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()); } 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); + }); + }); }); 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!`); }