Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
31 commits
Select commit Hold shift + click to select a range
0f0bdca
increase polling frequency and limit
stefangutica Aug 21, 2026
bc04011
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Aug 21, 2026
3e210d2
fixes + improvements
stefangutica Aug 26, 2026
5142e3a
skip canBeIgnored flagged transfers
stefangutica Aug 27, 2026
2ba6082
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Aug 28, 2026
fa8bb87
Merge remote-tracking branch 'origin/websocket-subscriptions-improvem…
stefangutica Aug 28, 2026
d1f5d96
reduce ws subscriptions broadcast intervbal
stefangutica Aug 28, 2026
268413f
performance improvements
stefangutica Aug 28, 2026
f460ef3
delete log
stefangutica Aug 28, 2026
bd74011
Merge pull request #1626 from multiversx/ws-fixes
stefangutica Aug 28, 2026
da124c2
add WS messages compression
stefangutica Aug 28, 2026
f0d3891
remove comm
stefangutica Aug 31, 2026
44b85b3
don't use cache on stats fetch
stefangutica Sep 1, 2026
a4e875f
use deterministic payload stringify on subscribe
stefangutica Sep 1, 2026
b43257c
use deterministic substituion instead of raw replace
stefangutica Sep 1, 2026
f2d876a
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Sep 3, 2026
ce6cb3f
disable tx pool
stefangutica Sep 7, 2026
8b687c2
disable pool in warmer
stefangutica Sep 7, 2026
2fdf4a7
enable tx pool
stefangutica Sep 7, 2026
8083b18
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Sep 7, 2026
7bea7ee
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Sep 8, 2026
1681782
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Sep 8, 2026
fcf1078
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Sep 10, 2026
920cfc3
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Sep 10, 2026
44a72b5
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Sep 11, 2026
642484b
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Sep 23, 2026
05cbffc
exclude sockets own rooms from global subscriptions limit
stefangutica Sep 24, 2026
f6fad85
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Sep 24, 2026
a5af502
Merge pull request #1641 from multiversx/ws-global-subscriptions-limit
stefangutica Sep 24, 2026
978f2ef
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Sep 24, 2026
bea6829
Merge remote-tracking branch 'origin/development' into websocket-subs…
stefangutica Sep 29, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion config/config.devnet.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,8 @@ features:
port: 6002
maxSubscriptionsPerInstance: 10000
maxSubscriptionsPerClient: 10
broadcastIntervalMs: 1000
broadcastIntervalMs: 600
compressionThreshold: 1024
eventsNotifier:
enabled: false
port: 5674
Expand Down
1 change: 1 addition & 0 deletions config/config.e2e-mocked.mainnet.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ features:
maxSubscriptionsPerInstance: 10000
maxSubscriptionsPerClient: 10
broadcastIntervalMs: 600
compressionThreshold: 1024
dataApi:
enabled: false
serviceUrl: 'https://data-api.multiversx.com'
Expand Down
1 change: 1 addition & 0 deletions config/config.e2e.mainnet.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ features:
maxSubscriptionsPerInstance: 10000
maxSubscriptionsPerClient: 10
broadcastIntervalMs: 600
compressionThreshold: 1024
eventsNotifier:
enabled: false
port: 5674
Expand Down
1 change: 1 addition & 0 deletions config/config.mainnet.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ features:
maxSubscriptionsPerInstance: 10000
maxSubscriptionsPerClient: 10
broadcastIntervalMs: 600
compressionThreshold: 1024
eventsNotifier:
enabled: false
port: 5674
Expand Down
3 changes: 2 additions & 1 deletion config/config.testnet.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,8 @@ features:
port: 6002
maxSubscriptionsPerInstance: 10000
maxSubscriptionsPerClient: 10
broadcastIntervalMs: 1000
broadcastIntervalMs: 600
compressionThreshold: 1024
eventsNotifier:
enabled: false
port: 5674
Expand Down
4 changes: 4 additions & 0 deletions src/common/api-config/api.config.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1016,6 +1016,10 @@ export class ApiConfigService {
return this.configService.get<number>('features.websocketSubscription.broadcastIntervalMs') ?? 1000;
}

getWebsocketSubscriptionCompressionThreshold(): number | undefined {
return this.configService.get<number>('features.websocketSubscription.compressionThreshold');
}

getWebsocketMaxSubscriptionsPerInstance(): number {
return this.configService.get<number>('features.websocketSubscription.maxSubscriptionsPerInstance') ?? 10_000;
}
Expand Down
24 changes: 24 additions & 0 deletions src/common/websockets/subscription-socket-adapter.ts
Original file line number Diff line number Diff line change
@@ -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 },
});
}
}
5 changes: 3 additions & 2 deletions src/crons/websocket/blocks.gateway.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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)) {
Expand All @@ -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)) {
Expand Down
128 changes: 128 additions & 0 deletions src/crons/websocket/custom.subscriptions.data.fetcher.ts
Original file line number Diff line number Diff line change
@@ -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<CustomSubscriptionsRoundData>) {
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<CustomSubscriptionsRoundData> {
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<Transaction[]> {
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<Events[]> {
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;
}
}
29 changes: 2 additions & 27 deletions src/crons/websocket/events.custom.gateway.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';

Expand All @@ -22,10 +19,6 @@ export class EventsCustomGateway {
@WebSocketServer()
server!: Server;

constructor(
private readonly eventsService: EventsService,
) { }

@UseInterceptors(LockingGuardInterceptor)
@SubscribeMessage('subscribeCustomEvents')
async handleCustomSubscription(
Expand Down Expand Up @@ -57,29 +50,11 @@ export class EventsCustomGateway {
return { status: 'unsubscribed' };
}

async pushEventsForTimestampMs(timestampMs: number): Promise<void> {
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<string, Events[]> = new Map();

for (const event of allEvents) {
for (const event of events) {
const roomKeys = RoomKeyGenerator.generate(
EventsCustomGateway.keyPrefix,
event,
Expand Down
2 changes: 1 addition & 1 deletion src/crons/websocket/events.gateway.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
2 changes: 1 addition & 1 deletion src/crons/websocket/network.gateway.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
11 changes: 11 additions & 0 deletions src/crons/websocket/room.key.generator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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, any>): string {
return JSON.stringify(
Object.keys(obj)
Expand Down
29 changes: 2 additions & 27 deletions src/crons/websocket/transaction.custom.gateway.ts
Original file line number Diff line number Diff line change
@@ -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';
Expand All @@ -20,10 +17,6 @@ export class TransactionsCustomGateway {
@WebSocketServer()
server!: Server;

constructor(
private readonly transactionService: TransactionService,
) { }

@UseInterceptors(LockingGuardInterceptor)
@SubscribeMessage('subscribeCustomTransactions')
async handleCustomSubscription(
Expand Down Expand Up @@ -52,28 +45,10 @@ export class TransactionsCustomGateway {
return { status: 'unsubscribed' };
}

async pushTransactionsForTimestampMs(timestampMs: number): Promise<void> {
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<string, Transaction[]> = new Map();
for (const transaction of allTransactions) {
for (const transaction of transactions) {
const roomKeys = RoomKeyGenerator.generate(
TransactionsCustomGateway.keyPrefix,
transaction,
Expand Down
2 changes: 1 addition & 1 deletion src/crons/websocket/transaction.gateway.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)) {
Expand Down
Loading
Loading