Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
152027e
take pool count from gateway and handle pool too large
stefangutica Sep 29, 2026
15f5290
warm pool counts in cache warmer and drop gateway count fallback
stefangutica Sep 29, 2026
099d690
cache pool too large as null in txpool key instead of separate key
stefangutica Sep 29, 2026
d51c316
explicit cacheNullable for txpool and warm pool counts in parallel
stefangutica Sep 29, 2026
9442566
answer pool too large with custom response on rest and websocket
stefangutica Sep 29, 2026
3dcaea9
send pool status on websocket and warm pool counts with promise all
stefangutica Sep 29, 2026
5c203d3
count pool by length and use gateway count only when pool is too large
stefangutica Sep 29, 2026
5d0aad1
send null pool on websocket and count types through gateway when pool…
stefangutica Sep 29, 2026
ec134c4
pass type to gateway pool count
stefangutica Sep 29, 2026
59e6fc6
back to 503 for pool too large, status on pool websocket and cache in…
stefangutica Sep 29, 2026
fb11296
count total from gateway only when pool is too large, without type
stefangutica Sep 29, 2026
f3a207b
use loose null checks for pool
stefangutica Sep 29, 2026
58307a9
return total pool count for type filter when pool is too large
stefangutica Sep 29, 2026
68166cb
return total pool count for any filter when pool is too large
stefangutica Sep 29, 2026
84a7933
use tx pool warmer ttl from config in cache warmer
stefangutica Sep 29, 2026
144591b
return null from gateway when pool is too large
stefangutica Sep 29, 2026
d100e98
send internal server error status on pool websocket
stefangutica Sep 30, 2026
29ab1df
count pool through gateway whenever pool fails
stefangutica Sep 30, 2026
24bf10f
use null pool with status for any pool error on websocket
stefangutica Sep 30, 2026
f4332d3
keep price per unit as before and fail pool with filters when pool is…
stefangutica Sep 30, 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
18 changes: 16 additions & 2 deletions src/common/gateway/gateway.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -182,8 +182,22 @@ export class GatewayService {
return new NftData(result.tokenData);
}

async getTransactionPool(): Promise<TxPoolGatewayResponse> {
return await this.get(`transaction/pool?fields=*`, GatewayComponentRequest.transactionPool);
async getTransactionPool(): Promise<TxPoolGatewayResponse | null> {
try {
return await this.get(`transaction/pool?fields=*`, GatewayComponentRequest.transactionPool);
} catch (error: any) {
if (error?.message?.startsWith('maxContentLength size of')) {
return null;
}

throw error;
}
}

async getTransactionPoolCount(): Promise<number> {
const result = await this.get('transaction/pool/count', GatewayComponentRequest.transactionPool);

return Object.values<number>(result.txPoolCounts).reduce((total, count) => total + count, 0);
}

async getTransaction(txHash: string): Promise<Transaction | undefined> {
Expand Down
12 changes: 11 additions & 1 deletion src/crons/cache.warmer/cache.warmer.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -151,9 +151,19 @@ export class CacheWarmerService {

@Lock({ name: 'Transaction pool invalidation', verbose: true })
async handleTxPoolInvalidations() {
const ttl = this.apiConfigService.getTransactionPoolCacheWarmerTtlInSeconds();
const pool = await this.poolService.getTxPoolRaw();

await this.invalidateKey(CacheInfo.TransactionPool.key, pool, this.apiConfigService.getTransactionPoolCacheWarmerTtlInSeconds());
const invalidations = [
this.invalidateKey(CacheInfo.TransactionPool.key, pool, ttl),
];

if (pool == null) {
const count = await this.gatewayService.getTransactionPoolCount();
invalidations.push(this.invalidateKey(CacheInfo.TransactionPoolCount.key, count, ttl));
}

await Promise.all(invalidations);
}

@Cron('*/2 * * * *')
Expand Down
16 changes: 14 additions & 2 deletions src/crons/websocket/pool.gateway.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import { QueryPagination } from 'src/common/entities/query.pagination';
import { PoolSubscribePayload } from '../../endpoints/pool/entities/pool.subscribe';
import { RoomKeyGenerator } from './room.key.generator';
import { LockingGuardInterceptor } from 'src/utils/locking.guard.interceptor';
import { PoolUpdateStatus } from '../../endpoints/pool/entities/pool.update.status';

@UseFilters(WebsocketExceptionsFilter)
@WebSocketGateway({ cors: { origin: '*' }, path: '/ws/subscription' })
Expand Down Expand Up @@ -63,20 +64,31 @@ export class PoolGateway {
type: filter.type,
});

let status = PoolUpdateStatus.success;

const [pool, poolCount] = await Promise.all([
this.poolService.getPool(
new QueryPagination({
from: filter.from,
size: filter.size,
}),
poolFilter,
),
).catch((error) => {
this.logger.error(error);
status = PoolUpdateStatus.internalServerError;
return null;
}),
this.poolService.getPoolCount(poolFilter),
]);

this.server.to(roomName).emit("poolUpdate", { pool, poolCount });
if (pool == null && status === PoolUpdateStatus.success) {
status = PoolUpdateStatus.tooLarge;
}

this.server.to(roomName).emit("poolUpdate", { status, pool, poolCount });
} catch (error) {
this.logger.error(error);
this.server.to(roomName).emit("poolUpdate", { status: PoolUpdateStatus.internalServerError, pool: null, poolCount: null });
}
}

Expand Down
5 changes: 5 additions & 0 deletions src/endpoints/pool/entities/pool.update.status.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
export enum PoolUpdateStatus {
success = 'success',
tooLarge = 'tooLarge',
internalServerError = 'internalServerError',
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
import { HttpException, HttpStatus } from "@nestjs/common";

export class TransactionPoolTooLargeException extends HttpException {
static readonly code = 'transaction_pool_too_large';

constructor() {
super({
statusCode: HttpStatus.SERVICE_UNAVAILABLE,
code: TransactionPoolTooLargeException.code,
message: 'The transaction pool is too large to be displayed',
}, HttpStatus.SERVICE_UNAVAILABLE);
}
}
17 changes: 15 additions & 2 deletions src/endpoints/pool/pool.controller.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,13 @@
import { ParseAddressAndMetachainPipe, ParseAddressPipe, ParseEnumPipe, ParseIntPipe, ParseTransactionHashPipe, ParseArrayPipe } from "@multiversx/sdk-nestjs-common";
import { Controller, DefaultValuePipe, Get, NotFoundException, Param, Query } from "@nestjs/common";
import { ApiExcludeEndpoint, ApiNotFoundResponse, ApiOkResponse, ApiOperation, ApiQuery, ApiTags } from "@nestjs/swagger";
import { ApiExcludeEndpoint, ApiNotFoundResponse, ApiOkResponse, ApiOperation, ApiQuery, ApiServiceUnavailableResponse, ApiTags } from "@nestjs/swagger";
import { PoolService } from "./pool.service";
import { QueryPagination } from "src/common/entities/query.pagination";
import { TransactionInPool } from "./entities/transaction.in.pool.dto";
import { TransactionType } from "../transactions/entities/transaction.type";
import { PoolFilter } from "./entities/pool.filter";
import { ParseArrayPipeOptions } from "@multiversx/sdk-nestjs-common/lib/pipes/entities/parse.array.options";
import { TransactionPoolTooLargeException } from "./entities/transaction.pool.too.large.exception";

@Controller()
@ApiTags('pool')
Expand All @@ -18,6 +19,7 @@ export class PoolController {
@Get("/pool")
@ApiOperation({ summary: 'Transactions pool', description: 'Returns the transactions that are currently in the memory pool.' })
@ApiOkResponse({ type: TransactionInPool, isArray: true })
@ApiServiceUnavailableResponse({ description: 'The transaction pool is too large to be displayed' })
@ApiQuery({ name: 'from', description: 'Number of items to skip for the result set', required: false })
@ApiQuery({ name: 'size', description: 'Number of items to retrieve', required: false })
@ApiQuery({ name: 'sender', description: 'Search in transaction pool by a specific sender', required: false })
Expand All @@ -37,14 +39,20 @@ export class PoolController {
@Query('type', new ParseEnumPipe(TransactionType)) type?: TransactionType,
@Query('function', new ParseArrayPipe(new ParseArrayPipeOptions({ allowEmptyString: true }))) functions?: string[],
): Promise<TransactionInPool[]> {
return await this.poolService.getPool(new QueryPagination({ from, size }), new PoolFilter({
const pool = await this.poolService.getPool(new QueryPagination({ from, size }), new PoolFilter({
sender: sender,
receiver: receiver,
senderShard: senderShard,
receiverShard: receiverShard,
type: type,
functions: functions,
}));

if (pool == null) {
throw new TransactionPoolTooLargeException();
}

return pool;
}

@Get("/pool/count")
Expand Down Expand Up @@ -85,10 +93,15 @@ export class PoolController {
@ApiOperation({ summary: 'Transaction from pool', description: 'Returns a transaction from the memory pool.' })
@ApiOkResponse({ type: TransactionInPool })
@ApiNotFoundResponse({ description: 'Transaction not found' })
@ApiServiceUnavailableResponse({ description: 'The transaction pool is too large to be displayed' })
async getTransactionFromPool(
@Param('txhash', ParseTransactionHashPipe) txHash: string,
): Promise<TransactionInPool> {
const transaction = await this.poolService.getTransactionFromPool(txHash);
if (transaction === null) {
throw new TransactionPoolTooLargeException();
}

if (transaction === undefined) {
throw new NotFoundException('Transaction not found');
}
Expand Down
52 changes: 41 additions & 11 deletions src/endpoints/pool/pool.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import { TxInPoolFields } from "src/common/gateway/entities/tx.in.pool.fields";
import { TransactionActionService } from "../transactions/transaction-action/transaction.action.service";
import { Transaction } from "../transactions/entities/transaction";
import { ApiUtils } from "@multiversx/sdk-nestjs-http";
import { TransactionPoolTooLargeException } from "./entities/transaction.pool.too.large.exception";

@Injectable()
export class PoolService {
Expand All @@ -22,43 +23,72 @@ export class PoolService {
private readonly transactionActionService: TransactionActionService,
) { }

async getTransactionFromPool(txHash: string): Promise<TransactionInPool | undefined> {
const pool = await this.getPoolWithFilters();
async getTransactionFromPool(txHash: string): Promise<TransactionInPool | undefined | null> {
const pool = await this.getTxPool();
if (pool == null) {
return null;
}

return pool.find(tx => tx.txHash === txHash);
}

async getPoolCount(filter: PoolFilter): Promise<number> {
const pool = await this.getPoolWithFilters(filter);
return pool.length;
const pool = await this.getTxPool().catch(() => null);
if (pool != null) {
return this.applyFilters(pool, filter).length;
}

return await this.cacheService.getOrSet(
CacheInfo.TransactionPoolCount.key,
async () => await this.gatewayService.getTransactionPoolCount(),
CacheInfo.TransactionPoolCount.ttl,
);
}

async getPool(
queryPagination: QueryPagination,
filter?: PoolFilter,
): Promise<TransactionInPool[]> {
): Promise<TransactionInPool[] | null> {
if (!this.apiConfigService.isTransactionPoolEnabled()) {
return [];
}

const { from, size } = queryPagination;
const pool = await this.getPoolWithFilters(filter);
return pool.slice(from, from + size);
const pool = await this.getTxPool();
if (pool == null) {
return null;
}

return this.applyFilters(pool, filter).slice(from, from + size);
}

async getPoolWithFilters(
filter?: PoolFilter,
): Promise<TransactionInPool[]> {
const pool = await this.cacheService.getOrSet(
const pool = await this.getTxPool();
if (pool == null) {
throw new TransactionPoolTooLargeException();
}

return this.applyFilters(pool, filter);
}

private async getTxPool(): Promise<TransactionInPool[] | null> {
return await this.cacheService.getOrSet(
CacheInfo.TransactionPool.key,
async () => await this.getTxPoolRaw(),
CacheInfo.TransactionPool.ttl,
CacheInfo.TransactionPool.ttl,
true,
);

return this.applyFilters(pool, filter);
}

async getTxPoolRaw(): Promise<TransactionInPool[]> {
async getTxPoolRaw(): Promise<TransactionInPool[] | null> {
const pool = await this.gatewayService.getTransactionPool();
if (pool == null) {
return null;
}

return this.parseTransactions(pool);
}

Expand Down
51 changes: 51 additions & 0 deletions src/test/unit/services/cache.warmer.pool.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
import { Locker, LockResult } from "@multiversx/sdk-nestjs-common";
import { CacheWarmerService } from "src/crons/cache.warmer/cache.warmer.service";
import { TransactionType } from "src/endpoints/transactions/entities/transaction.type";
import { CacheInfo } from "src/utils/cache.info";

describe('CacheWarmerService transaction pool', () => {
const ttl = 10;
const pool = [{ txHash: 'a', type: TransactionType.Transaction }];

let warmer: CacheWarmerService;
let poolService: any;
let gatewayService: any;
let cachingService: any;

beforeEach(() => {
jest.spyOn(Locker, 'lock').mockImplementation(async (_key: string, func: () => Promise<void>) => {
await func();
return LockResult.SUCCESS;
});

poolService = { getTxPoolRaw: jest.fn().mockResolvedValue(pool) };
gatewayService = { getTransactionPoolCount: jest.fn().mockResolvedValue(12) };
cachingService = { set: jest.fn() };

warmer = Object.assign(Object.create(CacheWarmerService.prototype), {
poolService,
gatewayService,
cachingService,
apiConfigService: { getTransactionPoolCacheWarmerTtlInSeconds: () => ttl },
clientProxy: { emit: jest.fn() },
});
});

it('should warm only the pool when it can be read', async () => {
await warmer.handleTxPoolInvalidations();

expect(cachingService.set).toHaveBeenCalledTimes(1);
expect(cachingService.set).toHaveBeenCalledWith(CacheInfo.TransactionPool.key, pool, ttl);
expect(gatewayService.getTransactionPoolCount).not.toHaveBeenCalled();
});

it('should warm the pool as null and refresh the count from the gateway when the pool is too large', async () => {
poolService.getTxPoolRaw.mockResolvedValue(null);

await warmer.handleTxPoolInvalidations();

expect(cachingService.set).toHaveBeenCalledTimes(2);
expect(cachingService.set).toHaveBeenCalledWith(CacheInfo.TransactionPool.key, null, ttl);
expect(cachingService.set).toHaveBeenCalledWith(CacheInfo.TransactionPoolCount.key, 12, ttl);
});
});
35 changes: 35 additions & 0 deletions src/test/unit/services/gateway.pool.count.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
import { GatewayService } from "src/common/gateway/gateway.service";

describe('GatewayService transaction pool', () => {
let gatewayService: GatewayService;
let apiService: any;

beforeEach(() => {
apiService = {
get: jest.fn().mockResolvedValue({ data: { data: { txPoolCounts: { '0': 1, '1': 2, '2': 3, '4294967295': 4 } } } }),
};

gatewayService = new GatewayService(
{ getGatewayUrl: () => 'https://gateway', getSnapshotlessGatewayUrl: () => undefined } as any,
apiService,
);
Object.assign(gatewayService, { eventEmitter: { emit: jest.fn() } });
});

it('should read a pool over the response size limit as null', async () => {
apiService.get.mockRejectedValue({ message: 'maxContentLength size of 2097152 exceeded' });

expect(await gatewayService.getTransactionPool()).toBeNull();
});

it('should keep throwing other failures of the pool', async () => {
apiService.get.mockRejectedValue({ message: 'connect ECONNREFUSED' });

await expect(gatewayService.getTransactionPool()).rejects.toEqual({ message: 'connect ECONNREFUSED' });
});

it('should sum the counts of every shard', async () => {
expect(await gatewayService.getTransactionPoolCount()).toStrictEqual(10);
expect(apiService.get).toHaveBeenCalledWith('https://gateway/transaction/pool/count', expect.anything(), undefined);
});
});
Loading
Loading